pub struct Service {
pub db: Data,
server: Arc<Server>,
services: Arc<OnceServices>,
channels: Vec<(Sender<Msg>, Receiver<Msg>)>,
flushes: Mutex<JoinSet<()>>,
stalled: Mutex<HashMap<OwnedServerName, Option<Instant>>>,
}Expand description
Outbound delivery of PDUs and EDUs to federation peers, appservices, and push gateways.
Requests are written as durable queue rows and dispatched to a pool of sender workers sharded by destination.
Fields§
§db: Data§server: Arc<Server>§services: Arc<OnceServices>§channels: Vec<(Sender<Msg>, Receiver<Msg>)>§flushes: Mutex<JoinSet<()>>§stalled: Mutex<HashMap<OwnedServerName, Option<Instant>>>Implementations§
Source§impl Service
impl Service
Sourcepub async fn send_to_device_appservices<'a, I>(
&self,
sender: &'a UserId,
target_user: &'a UserId,
deliveries: I,
event_type: &'a str,
content: &'a Value,
) -> Result
pub async fn send_to_device_appservices<'a, I>( &self, sender: &'a UserId, target_user: &'a UserId, deliveries: I, event_type: &'a str, content: &'a Value, ) -> Result
Queue stored to-device events for delivery to interested appservices (MSC4203).
deliveries are the concrete recipient devices already written to the
inbox, after AllDevices expansion. Each delivery becomes one queue row
per interested appservice, written under one cork; with no interested
appservice nothing is serialized.
Source§impl Service
impl Service
Sourcepub async fn send_device_list_appservices(
&self,
user_id: &UserId,
count: u64,
) -> Result
pub async fn send_device_list_appservices( &self, user_id: &UserId, count: u64, ) -> Result
Queue a device_lists.changed marker (MSC3202) for delivery to
appservices that opted into transaction extensions and are interested
in user_id.
The caller passes the count it already allocated so the marker uniquifies the transaction hash.
Source§impl Service
impl Service
Whether user_id shares a device-list-interesting room with info.
A joined room the appservice participates in counts when it is encrypted,
or unconditionally when device_key_update_encrypted_rooms_only is off.
Source§impl Service
impl Service
pub(super) async fn send_events_dest_appservice( &self, id: String, events: Vec<SendingEvent>, ) -> Result<Destination, (Destination, Error)>
Source§impl Service
impl Service
async fn txn_part( &self, info: &RegistrationInfo, msc3202: bool, event: &SendingEvent, ) -> Part
Source§impl Service
impl Service
Sourceasync fn msc3202_key_counts(
&self,
users: BTreeSet<OwnedUserId>,
recipients: SmallVec<[(OwnedUserId, OwnedDeviceId); 1]>,
) -> (BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, BTreeMap<OneTimeKeyAlgorithm, UInt>>>, BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, Vec<OneTimeKeyAlgorithm>>>)
async fn msc3202_key_counts( &self, users: BTreeSet<OwnedUserId>, recipients: SmallVec<[(OwnedUserId, OwnedDeviceId); 1]>, ) -> (BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, BTreeMap<OneTimeKeyAlgorithm, UInt>>>, BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, Vec<OneTimeKeyAlgorithm>>>)
MSC3202 one-time-key counts and unused fallback key types over every
device of users plus the specific recipients.
Recomputed per build rather than snapshotted, so a retry ships fresh counts.
Source§impl Service
impl Service
Sourcepub(super) async fn send_events_dest_federation(
&self,
server: OwnedServerName,
events: Vec<SendingEvent>,
) -> (Result<Destination, (Destination, Error)>, bool)
pub(super) async fn send_events_dest_federation( &self, server: OwnedServerName, events: Vec<SendingEvent>, ) -> (Result<Destination, (Destination, Error)>, bool)
Send a federation transaction, reporting whether one went out at all.
Rows that all fail to load leave nothing to send; they still succeed, so their keys are acknowledged.
Source§impl Service
impl Service
Sourcepub fn schedule_flush_suppressed_for_pushkey(
&self,
user_id: OwnedUserId,
pushkey: String,
reason: &'static str,
)
pub fn schedule_flush_suppressed_for_pushkey( &self, user_id: OwnedUserId, pushkey: String, reason: &'static str, )
Schedule a flush of the pushes suppressed for one pushkey.
The flush runs as a task this service owns, so the caller never waits on the push gateway.
Source§impl Service
impl Service
Sourcepub fn schedule_flush_suppressed_for_user(
&self,
user_id: OwnedUserId,
reason: &'static str,
)
pub fn schedule_flush_suppressed_for_user( &self, user_id: OwnedUserId, reason: &'static str, )
Schedule a flush of the pushes suppressed for every pushkey a user owns.
The flush runs as a task this service owns, so the caller never waits on the push gateway.
Source§impl Service
impl Service
async fn flush_suppressed_for_pushkey( &self, user_id: &UserId, pushkey: &str, reason: &'static str, )
Source§impl Service
impl Service
async fn flush_suppressed_for_user( &self, user_id: &UserId, reason: &'static str, )
Source§impl Service
impl Service
async fn flush_suppressed_rooms( &self, flush: &Flush<'_>, rooms: SuppressedRooms, )
Source§impl Service
impl Service
pub(super) async fn enqueue_suppressed_push_events( &self, user_id: &UserId, pushkey: &str, events: &[SendingEvent], ) -> usize
Source§impl Service
impl Service
Sourceasync fn suppressable_pdu(
&self,
user_id: &UserId,
pdu_id: &RawPduId,
) -> Option<Pdu>
async fn suppressable_pdu( &self, user_id: &UserId, pdu_id: &RawPduId, ) -> Option<Pdu>
Load a suppressed PDU if a push for it is still worth sending.
A missing or redacted PDU is dropped with a log line and never notified.
Source§impl Service
impl Service
Sourcepub(super) async fn pushing_suppressed(&self, user_id: &UserId) -> bool
pub(super) async fn pushing_suppressed(&self, user_id: &UserId) -> bool
Decide whether pushes for a user are suppressed as active.
The heuristic combines the presence age and the most recent sync gap, and
only applies when suppress_push_when_active is enabled. An offline user
is never suppressed; the two ACTIVE_* constants are the thresholds.
Source§impl Service
impl Service
pub(super) async fn send_events_dest_push( &self, user_id: OwnedUserId, pushkey: String, events: Vec<SendingEvent>, ) -> Result<Destination, (Destination, Error)>
Source§impl Service
impl Service
pub(super) fn send_events( &self, dest: Destination, items: Vec<(Vec<u8>, SendingEvent)>, split: Option<Split>, ) -> BoxFuture<'_, Completion>
Source§impl Service
impl Service
pub(super) async fn startup_netburst<'a>( &'a self, id: usize, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
fn mark_pending(&self, dest: &Destination)
Source§impl Service
impl Service
async fn arm_inherited_destinations( &self, id: usize, statuses: &HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
pub(super) async fn arm_startup_wake( &self, dest: Destination, index: u64, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, ) -> u64
Source§impl Service
impl Service
pub(super) async fn handle_response<'a>( &'a self, __arg1: Completion, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
async fn handle_response_ok<'a>( &'a self, dest: Destination, split: Option<Split>, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
async fn handle_response_err<'a>( &'a self, __arg1: (Destination, Error), split: Option<Split>, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
pub(super) async fn handle_force_retry<'a>( &'a self, dest: Destination, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
Sourcepub(super) async fn select_edus_device_changes(
&self,
server_name: &ServerName,
since: (u64, u64),
max_edu_count: &AtomicU64,
events_len: &AtomicUsize,
) -> Selected
pub(super) async fn select_edus_device_changes( &self, server_name: &ServerName, since: (u64, u64), max_edu_count: &AtomicU64, events_len: &AtomicUsize, ) -> Selected
Select device-list deltas and signing-key updates for local users.
Legacy records and windows exceeding the device limit force a snapshot resync.
Source§impl Service
impl Service
async fn select_user_devices( &self, selected: Selected, user_id: &UserId, counts: BTreeSet<u64>, server_name: &ServerName, events_len: &AtomicUsize, ) -> Selected
Source§impl Service
impl Service
async fn signing_key_edu( &self, user_id: &UserId, server_name: &ServerName, ) -> Option<EduBuf>
Source§impl Service
impl Service
async fn device_delta_edu(&self, user_id: &UserId, delta: Delta<'_>) -> EduBuf
Source§impl Service
impl Service
Sourcepub(super) async fn select_edus_presence(
&self,
server_name: &ServerName,
since: (u64, u64),
max_edu_count: &AtomicU64,
events_len: &AtomicUsize,
) -> Option<EduBuf>
pub(super) async fn select_edus_presence( &self, server_name: &ServerName, since: (u64, u64), max_edu_count: &AtomicU64, events_len: &AtomicUsize, ) -> Option<EduBuf>
Select one presence EDU for the users a server may see.
The window’s local presence transitions collapse into one Edu::Presence
carrying the latest state per user. A budget trip drops it, and presence
self-heals on the next transition.
Source§impl Service
impl Service
Sourceasync fn presence_update(
&self,
server_name: &ServerName,
user_id: &UserId,
presence_bytes: &[u8],
) -> Option<(OwnedUserId, PresenceUpdate)>
async fn presence_update( &self, server_name: &ServerName, user_id: &UserId, presence_bytes: &[u8], ) -> Option<(OwnedUserId, PresenceUpdate)>
The presence update to ship for one transition, if the server may see it.
Source§impl Service
impl Service
Sourcepub(super) async fn select_edus_receipts(
&self,
server_name: &ServerName,
since: (u64, u64),
max_edu_count: &AtomicU64,
events_len: &AtomicUsize,
) -> Selected
pub(super) async fn select_edus_receipts( &self, server_name: &ServerName, since: (u64, u64), max_edu_count: &AtomicU64, events_len: &AtomicUsize, ) -> Selected
Select read-receipt EDUs across every room shared with the server.
MSC3771 lets a user emit multiple receipts in the same EDU window, one
per thread context. The federation EDU shape allows only one
ReceiptData per (room, user) slot, so a user with N parallel
thread receipts ships across N parallel Edu::Receipt buffers within
the same transaction. Each buffer is shape-compliant; receivers
process them as independent receipt EDUs and our storage keeps each
thread distinct.
Source§impl Service
impl Service
Sourceasync fn select_edus_receipts_room(
&self,
room_id: &RoomId,
since: (u64, u64),
max_edu_count: &AtomicU64,
num: &AtomicUsize,
) -> SmallVec<[ReceiptMap; 1]> ⓘ
async fn select_edus_receipts_room( &self, room_id: &RoomId, since: (u64, u64), max_edu_count: &AtomicU64, num: &AtomicUsize, ) -> SmallVec<[ReceiptMap; 1]> ⓘ
Look for read receipts in this room.
The receipt-limit budget bounds distinct users only; subsequent thread receipts for an already-counted user do not consume additional budget.
Source§impl Service
impl Service
pub(super) async fn select_events( &self, dest: &Destination, new_events: SmallVec<[(Vec<u8>, SendingEvent); 1]>, statuses: &mut HashMap<Destination, TransactionStatus>, ) -> Result<Selection>
Source§impl Service
impl Service
async fn select_events_current( &self, dest: &Destination, statuses: &mut HashMap<Destination, TransactionStatus>, retry_action: RetryAction, ) -> Current
Source§impl Service
impl Service
Sourcefn transition(
&self,
dest: &Destination,
status: &mut TransactionStatus,
retry_action: RetryAction,
) -> Current
fn transition( &self, dest: &Destination, status: &mut TransactionStatus, retry_action: RetryAction, ) -> Current
Advance a destination’s status for a new selection.
Distinguishes busy destinations from permitted selection and active replay.
Source§impl Service
impl Service
fn clear_stalled(&self, dest: &Destination)
Source§impl Service
impl Service
Sourcepub(super) async fn federation_batch(
&self,
dest: &Destination,
server: &ServerName,
new_events: SmallVec<[(Vec<u8>, SendingEvent); 1]>,
) -> Selection
pub(super) async fn federation_batch( &self, dest: &Destination, server: &ServerName, new_events: SmallVec<[(Vec<u8>, SendingEvent); 1]>, ) -> Selection
Compose a federation destination’s next transaction around its parked rooms.
Each expired park is retried alone first, and one whose room has nothing queued ends. Rows of other parked rooms stay queued, and a destination holding nothing else waits for the earliest park to expire.
Source§impl Service
impl Service
Sourceasync fn claim_new(
&self,
new_events: SmallVec<[(Vec<u8>, SendingEvent); 1]>,
skip: &[ShortRoomId],
) -> Vec<(Vec<u8>, SendingEvent)>
async fn claim_new( &self, new_events: SmallVec<[(Vec<u8>, SendingEvent); 1]>, skip: &[ShortRoomId], ) -> Vec<(Vec<u8>, SendingEvent)>
Claim a request’s own queue rows, leaving those of skipped rooms queued.
Flush markers carry no row and are dropped; the claimed rows turn active in one batch.
Source§impl Service
impl Service
Sourceasync fn with_edus(
&self,
dest: &Destination,
server_name: &ServerName,
items: Vec<(Vec<u8>, SendingEvent)>,
skip: &[ShortRoomId],
) -> Vec<(Vec<u8>, SendingEvent)>
async fn with_edus( &self, dest: &Destination, server_name: &ServerName, items: Vec<(Vec<u8>, SendingEvent)>, skip: &[ShortRoomId], ) -> Vec<(Vec<u8>, SendingEvent)>
Top up a federation transaction with the EDUs accrued since its last window.
An empty transaction first claims the head of the queue. Rows of the skipped rooms neither join the transaction nor hold back fresh EDUs.
Source§impl Service
impl Service
Sourcepub(super) async fn resume_queued(
&self,
dest: &Destination,
skip: &[ShortRoomId],
) -> Vec<(Vec<u8>, SendingEvent)>
pub(super) async fn resume_queued( &self, dest: &Destination, skip: &[ShortRoomId], ) -> Vec<(Vec<u8>, SendingEvent)>
Claim the head of a destination’s queue as its next transaction.
Rows of the skipped rooms are passed over; the claimed rows turn active in one batch.
Source§impl Service
impl Service
pub(super) async fn select_edus( &self, server_name: &ServerName, budget_used: usize, ) -> impl Iterator<Item = (Vec<u8>, SendingEvent)>
Source§impl Service
impl Service
Sourcepub(super) async fn split_failure(
&self,
server: &ServerName,
error: &Error,
split: Option<Split>,
tries: u32,
) -> (Option<Split>, u32)
pub(super) async fn split_failure( &self, server: &ServerName, error: &Error, split: Option<Split>, tries: u32, ) -> (Option<Split>, u32)
Choose how a failed federation transaction continues.
After repeated failures the transaction’s PDU rows return to the queue and its rooms go out one per transaction; a room failing after a delivery is parked. Returns the split to keep and the failure count its retry backs off by, zero when the split advanced.
Source§impl Service
impl Service
Sourceasync fn control(
&self,
dest: &Destination,
server: &ServerName,
rooms: &SmallVec<[ShortRoomId; 2]>,
) -> Option<ShortRoomId>
async fn control( &self, dest: &Destination, server: &ServerName, rooms: &SmallVec<[ShortRoomId; 2]>, ) -> Option<ShortRoomId>
Find the first queued room outside the split and the server’s parks.
A control delivered while the split’s rooms keep failing implicates those rooms rather than the destination.
Source§impl Service
impl Service
Sourcepub(super) async fn split_delivered(
&self,
dest: &Destination,
split: Split,
) -> Option<(Vec<(Vec<u8>, SendingEvent)>, Split)>
pub(super) async fn split_delivered( &self, dest: &Destination, split: Split, ) -> Option<(Vec<(Vec<u8>, SendingEvent)>, Split)>
Advance a split past its delivered head room, ending any park it had.
Returns the next room’s transaction, or nothing once the split is done.
Source§impl Service
impl Service
Sourcepub(super) async fn slice(
&self,
dest: &Destination,
split: Split,
) -> Option<(Vec<(Vec<u8>, SendingEvent)>, Split)>
pub(super) async fn slice( &self, dest: &Destination, split: Split, ) -> Option<(Vec<(Vec<u8>, SendingEvent)>, Split)>
Promote the head room’s queued rows as the split’s next transaction.
The split ends when it has no rooms left, or its head room has no rows. The rejected transaction’s EDU rows stay active and ride along until delivered.
Source§impl Service
impl Service
pub(super) async fn drain_due_wakes<'a>( &'a self, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
async fn handle_wake<'a>( &'a self, dest: Destination, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
async fn handle_federation_wake<'a>( &'a self, server: OwnedServerName, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
pub(super) async fn arm_federation_wake( &self, server: OwnedServerName, tries: u32, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
pub(super) fn arm_push_wake( &self, dest: Destination, error: &Error, statuses: &HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
pub(super) fn push_backoff_remaining( &self, status: Option<&TransactionStatus>, ) -> Option<Duration>
Source§impl Service
impl Service
async fn work_loop<'a>( &'a self, id: usize, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
async fn handle_request<'a>( &'a self, msg: Msg, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
fn schedule_events<'a>( &'a self, dest: Destination, selection: Selection, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, statuses: &mut HashMap<Destination, TransactionStatus>, wakes: &mut BinaryHeap<Reverse<(Instant, Destination)>>, )
Source§impl Service
impl Service
async fn finish_responses<'a>( &'a self, futures: &mut FuturesUnordered<BoxFuture<'a, Completion>>, )
Source§impl Service
impl Service
Sourcepub(super) async fn run(self: Arc<Self>) -> Result
pub(super) async fn run(self: Arc<Self>) -> Result
Run the sender workers to completion.
One worker is spawned per channel and joined; a panic among them is returned so the manager restarts the service. The suppressed-push flush tasks are then aborted and joined so a panic among them is still reported.
Source§impl Service
impl Service
pub(super) fn spawn_flush<F>(&self, flush: F)
Source§impl Service
impl Service
Sourcepub fn send_pdu_push(
&self,
pdu_id: &RawPduId,
user: &UserId,
pushkey: String,
) -> Result
pub fn send_pdu_push( &self, pdu_id: &RawPduId, user: &UserId, pushkey: String, ) -> Result
Queue a PDU for delivery to one of a user’s pushers.
The row is durable and the shard owning the destination is woken.
Source§impl Service
impl Service
fn queue_and_dispatch(&self, dest: Destination, event: SendingEvent) -> Result
Source§impl Service
impl Service
Sourcepub async fn refresh_push_badge(&self, user_id: &UserId) -> Result
pub async fn refresh_push_badge(&self, user_id: &UserId) -> Result
Queue a counts-only push refresh for every pusher owned by a user.
Rows are durable, coalesced, and recomputed at send time.
Source§impl Service
impl Service
Sourcepub fn send_pdu_appservice(
&self,
appservice_id: String,
pdu_id: RawPduId,
) -> Result
pub fn send_pdu_appservice( &self, appservice_id: String, pdu_id: RawPduId, ) -> Result
Queue a PDU for delivery to an appservice.
The row is durable and the shard owning the destination is woken.
Source§impl Service
impl Service
Sourcepub async fn send_pdu_room(&self, room_id: &RoomId, pdu_id: &RawPduId) -> Result
pub async fn send_pdu_room(&self, room_id: &RoomId, pdu_id: &RawPduId) -> Result
Queue a PDU for delivery to every remote server in a room.
The fan-out is one durable row per server, dispatched under a single cork.
Source§impl Service
impl Service
Sourcepub async fn send_pdu_servers<'a, S>(
&self,
servers: S,
pdu_id: &RawPduId,
) -> Resultwhere
S: Stream<Item = &'a ServerName> + Send + 'a,
pub async fn send_pdu_servers<'a, S>(
&self,
servers: S,
pdu_id: &RawPduId,
) -> Resultwhere
S: Stream<Item = &'a ServerName> + Send + 'a,
Queue a PDU for delivery to each of the given servers.
The fan-out is one durable row per server, dispatched under a single cork.
Source§impl Service
impl Service
async fn queue_and_dispatch_servers<'a, S>(
&self,
servers: S,
event: SendingEvent,
) -> Resultwhere
S: Stream<Item = &'a ServerName> + Send + 'a,
Source§impl Service
impl Service
Sourcepub fn send_edu_server(&self, server: &ServerName, serialized: EduBuf) -> Result
pub fn send_edu_server(&self, server: &ServerName, serialized: EduBuf) -> Result
Queue an EDU for delivery to a server.
The row is durable and the shard owning the destination is woken.
Source§impl Service
impl Service
Sourcepub async fn send_edu_room(
&self,
room_id: &RoomId,
serialized: EduBuf,
) -> Result
pub async fn send_edu_room( &self, room_id: &RoomId, serialized: EduBuf, ) -> Result
Queue an EDU for delivery to every remote server in a room.
The fan-out is one durable row per server, dispatched under a single cork.
Source§impl Service
impl Service
Sourcepub async fn send_edu_servers<'a, S>(
&self,
servers: S,
serialized: EduBuf,
) -> Resultwhere
S: Stream<Item = &'a ServerName> + Send + 'a,
pub async fn send_edu_servers<'a, S>(
&self,
servers: S,
serialized: EduBuf,
) -> Resultwhere
S: Stream<Item = &'a ServerName> + Send + 'a,
Queue an EDU for delivery to each of the given servers.
The fan-out is one durable row per server, dispatched under a single cork.
Source§impl Service
impl Service
Sourcepub async fn send_edu_room_appservices<'a, F>(
&self,
room_id: &RoomId,
serializer: F,
) -> Result
pub async fn send_edu_room_appservices<'a, F>( &self, room_id: &RoomId, serializer: F, ) -> Result
Queue an EDU for every appservice interested in a room.
An appservice is interested when it receives ephemeral events and the room
is in its namespace, it is present in the room, or one of the room’s local
aliases matches. The serializer writes EphemeralData, not a federation
Edu, once per matching appservice.
Source§impl Service
impl Service
Sourcepub fn send_edu_appservice(
&self,
appservice_id: String,
serialized: EduBuf,
) -> Result
pub fn send_edu_appservice( &self, appservice_id: String, serialized: EduBuf, ) -> Result
Queue an EDU for delivery to a specific appservice.
The row is durable and the shard owning the destination is woken.
Source§impl Service
impl Service
Sourcepub async fn flush_room(&self, room_id: &RoomId) -> Result
pub async fn flush_room(&self, room_id: &RoomId) -> Result
Wake the sender for every remote server in a room.
A flush is not queued as a row; it only prompts the shard to compose a transaction from whatever is pending.
Source§impl Service
impl Service
Sourcepub async fn flush_servers<'a, S>(&self, servers: S) -> Resultwhere
S: Stream<Item = &'a ServerName> + Send + 'a,
pub async fn flush_servers<'a, S>(&self, servers: S) -> Resultwhere
S: Stream<Item = &'a ServerName> + Send + 'a,
Wake the sender for each of the given servers.
A flush is not queued as a row; it only prompts the shard to compose a transaction from whatever is pending.
Source§impl Service
impl Service
fn dispatch_flush(&self, dest: Destination) -> Result
Source§impl Service
impl Service
Sourcepub fn flush_appservice(&self, appservice_id: String) -> Result
pub fn flush_appservice(&self, appservice_id: String) -> Result
Wake the sender for an appservice.
A flush is not queued as a row; it only prompts the shard to compose a transaction from whatever is pending.
Source§impl Service
impl Service
Sourcepub async fn notify_peer_alive(&self, server: &ServerName) -> bool
pub async fn notify_peer_alive(&self, server: &ServerName) -> bool
Wake the sender for a federation peer that has proven reachable.
Reachability comes from inbound activity or an operator reset. The flush resumes a waiting sender generation after its notification floor, or immediately when peer failure rows existed. The return value reports only whether those peer rows existed.
Source§impl Service
impl Service
Sourcepub async fn cleanup_events(
&self,
appservice_id: Option<&str>,
user_id: Option<&UserId>,
push_key: Option<&str>,
) -> Result
pub async fn cleanup_events( &self, appservice_id: Option<&str>, user_id: Option<&UserId>, push_key: Option<&str>, ) -> Result
Clean up queued sending event data.
Accepts either an appservice ID alone, after its registration is removed, or a user ID with a push key, after the pusher is deleted; any other combination is ignored with a warning.
Trait Implementations§
Source§impl Service for Service
impl Service for Service
Source§fn build(args: &Args<'_>) -> Result<Arc<Self>>
fn build(args: &Args<'_>) -> Result<Arc<Self>>
Source§fn worker<'async_trait>(
self: Arc<Self>,
) -> Pin<Box<dyn Future<Output = Result> + Send + 'async_trait>>where
Self: 'async_trait,
fn worker<'async_trait>(
self: Arc<Self>,
) -> Pin<Box<dyn Future<Output = Result> + Send + 'async_trait>>where
Self: 'async_trait,
Source§fn interrupt<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = ()> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn interrupt<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = ()> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Source§fn name(&self) -> &str
fn name(&self) -> &str
crate::service::make_name(std::module_path!())Source§fn unconstrained(&self) -> bool
fn unconstrained(&self) -> bool
Auto Trait Implementations§
impl !Freeze for Service
impl !RefUnwindSafe for Service
impl !UnwindSafe for Service
impl Send for Service
impl Sync for Service
impl Unpin for Service
impl UnsafeUnpin for Service
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> ExpectInto for T
impl<T> ExpectInto for T
Source§impl<T> Expected for T
impl<T> Expected for T
Source§fn expected_add(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedAdd,
fn expected_add(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedAdd,
rhs with an expectation that the operation is valid. Read moreSource§fn expected_sub(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedSub,
fn expected_sub(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedSub,
rhs with an expectation that the operation is valid. Read moreSource§fn expected_mul(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedMul,
fn expected_mul(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedMul,
rhs with an expectation that the operation is valid. Read moreSource§fn expected_div(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedDiv,
fn expected_div(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedDiv,
rhs with an expectation that the operation is valid. Read moreSource§fn expected_rem(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedRem,
fn expected_rem(self, rhs: Self) -> Selfwhere
Self: Sized + CheckedRem,
§impl<T> Identity for Twhere
T: ?Sized,
impl<T> Identity for Twhere
T: ?Sized,
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreimpl<T> JsonCastable<CanonicalJsonValue> for T
impl<T> JsonCastable<Value> for T
§impl<T> Paint for Twhere
T: ?Sized,
impl<T> Paint for Twhere
T: ?Sized,
§fn fg(&self, value: Color) -> Painted<&T>
fn fg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the foreground set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like red() and
green(), which have the same functionality but are
pithier.
§Example
Set foreground color to white using fg():
use yansi::{Paint, Color};
painted.fg(Color::White);Set foreground color to white using white().
use yansi::Paint;
painted.white();§fn bright_black(&self) -> Painted<&T>
fn bright_black(&self) -> Painted<&T>
§fn bright_red(&self) -> Painted<&T>
fn bright_red(&self) -> Painted<&T>
§fn bright_green(&self) -> Painted<&T>
fn bright_green(&self) -> Painted<&T>
§fn bright_yellow(&self) -> Painted<&T>
fn bright_yellow(&self) -> Painted<&T>
§fn bright_blue(&self) -> Painted<&T>
fn bright_blue(&self) -> Painted<&T>
§fn bright_magenta(&self) -> Painted<&T>
fn bright_magenta(&self) -> Painted<&T>
§fn bright_cyan(&self) -> Painted<&T>
fn bright_cyan(&self) -> Painted<&T>
§fn bright_white(&self) -> Painted<&T>
fn bright_white(&self) -> Painted<&T>
§fn bg(&self, value: Color) -> Painted<&T>
fn bg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the background set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like on_red() and
on_green(), which have the same functionality but
are pithier.
§Example
Set background color to red using fg():
use yansi::{Paint, Color};
painted.bg(Color::Red);Set background color to red using on_red().
use yansi::Paint;
painted.on_red();§fn on_primary(&self) -> Painted<&T>
fn on_primary(&self) -> Painted<&T>
§fn on_magenta(&self) -> Painted<&T>
fn on_magenta(&self) -> Painted<&T>
§fn on_bright_black(&self) -> Painted<&T>
fn on_bright_black(&self) -> Painted<&T>
§fn on_bright_red(&self) -> Painted<&T>
fn on_bright_red(&self) -> Painted<&T>
§fn on_bright_green(&self) -> Painted<&T>
fn on_bright_green(&self) -> Painted<&T>
§fn on_bright_yellow(&self) -> Painted<&T>
fn on_bright_yellow(&self) -> Painted<&T>
§fn on_bright_blue(&self) -> Painted<&T>
fn on_bright_blue(&self) -> Painted<&T>
§fn on_bright_magenta(&self) -> Painted<&T>
fn on_bright_magenta(&self) -> Painted<&T>
§fn on_bright_cyan(&self) -> Painted<&T>
fn on_bright_cyan(&self) -> Painted<&T>
§fn on_bright_white(&self) -> Painted<&T>
fn on_bright_white(&self) -> Painted<&T>
§fn attr(&self, value: Attribute) -> Painted<&T>
fn attr(&self, value: Attribute) -> Painted<&T>
Enables the styling [Attribute] value.
This method should be used rarely. Instead, prefer to use
attribute-specific builder methods like bold() and
underline(), which have the same functionality
but are pithier.
§Example
Make text bold using attr():
use yansi::{Paint, Attribute};
painted.attr(Attribute::Bold);Make text bold using using bold().
use yansi::Paint;
painted.bold();§fn rapid_blink(&self) -> Painted<&T>
fn rapid_blink(&self) -> Painted<&T>
§fn quirk(&self, value: Quirk) -> Painted<&T>
fn quirk(&self, value: Quirk) -> Painted<&T>
Enables the yansi [Quirk] value.
This method should be used rarely. Instead, prefer to use quirk-specific
builder methods like mask() and
wrap(), which have the same functionality but are
pithier.
§Example
Enable wrapping using .quirk():
use yansi::{Paint, Quirk};
painted.quirk(Quirk::Wrap);Enable wrapping using wrap().
use yansi::Paint;
painted.wrap();§fn clear(&self) -> Painted<&T>
👎Deprecated since 1.0.1: renamed to resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.
fn clear(&self) -> Painted<&T>
renamed to resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.
§fn whenever(&self, value: Condition) -> Painted<&T>
fn whenever(&self, value: Condition) -> Painted<&T>
Conditionally enable styling based on whether the [Condition] value
applies. Replaces any previous condition.
See the crate level docs for more details.
§Example
Enable styling painted only when both stdout and stderr are TTYs:
use yansi::{Paint, Condition};
painted.red().on_yellow().whenever(Condition::STDOUTERR_ARE_TTY);§impl<T> Pointable for T
impl<T> Pointable for T
§impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> PolicyExt for Twhere
T: ?Sized,
§impl<T> ServiceExt for T
impl<T> ServiceExt for T
§fn add_extension<T>(self, value: T) -> AddExtension<Self, T>where
Self: Sized,
fn add_extension<T>(self, value: T) -> AddExtension<Self, T>where
Self: Sized,
§fn compression(self) -> Compression<Self>where
Self: Sized,
fn compression(self) -> Compression<Self>where
Self: Sized,
§fn decompression(self) -> Decompression<Self>where
Self: Sized,
fn decompression(self) -> Decompression<Self>where
Self: Sized,
§fn trace_for_http(self) -> Trace<Self, SharedClassifier<ServerErrorsAsFailures>>where
Self: Sized,
fn trace_for_http(self) -> Trace<Self, SharedClassifier<ServerErrorsAsFailures>>where
Self: Sized,
§fn trace_for_grpc(self) -> Trace<Self, SharedClassifier<GrpcErrorsAsFailures>>where
Self: Sized,
fn trace_for_grpc(self) -> Trace<Self, SharedClassifier<GrpcErrorsAsFailures>>where
Self: Sized,
§fn follow_redirects(self) -> FollowRedirect<Self>where
Self: Sized,
fn follow_redirects(self) -> FollowRedirect<Self>where
Self: Sized,
§fn sensitive_headers(
self,
headers: impl IntoIterator<Item = HeaderName>,
) -> SetSensitiveRequestHeaders<SetSensitiveResponseHeaders<Self>>where
Self: Sized,
fn sensitive_headers(
self,
headers: impl IntoIterator<Item = HeaderName>,
) -> SetSensitiveRequestHeaders<SetSensitiveResponseHeaders<Self>>where
Self: Sized,
§fn sensitive_request_headers(
self,
headers: impl IntoIterator<Item = HeaderName>,
) -> SetSensitiveRequestHeaders<Self>where
Self: Sized,
fn sensitive_request_headers(
self,
headers: impl IntoIterator<Item = HeaderName>,
) -> SetSensitiveRequestHeaders<Self>where
Self: Sized,
§fn sensitive_response_headers(
self,
headers: impl IntoIterator<Item = HeaderName>,
) -> SetSensitiveResponseHeaders<Self>where
Self: Sized,
fn sensitive_response_headers(
self,
headers: impl IntoIterator<Item = HeaderName>,
) -> SetSensitiveResponseHeaders<Self>where
Self: Sized,
§fn override_request_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetRequestHeader<Self, M>where
Self: Sized,
fn override_request_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetRequestHeader<Self, M>where
Self: Sized,
§fn append_request_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetRequestHeader<Self, M>where
Self: Sized,
fn append_request_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetRequestHeader<Self, M>where
Self: Sized,
§fn insert_request_header_if_not_present<M>(
self,
header_name: HeaderName,
make: M,
) -> SetRequestHeader<Self, M>where
Self: Sized,
fn insert_request_header_if_not_present<M>(
self,
header_name: HeaderName,
make: M,
) -> SetRequestHeader<Self, M>where
Self: Sized,
§fn override_response_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetResponseHeader<Self, M>where
Self: Sized,
fn override_response_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetResponseHeader<Self, M>where
Self: Sized,
§fn append_response_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetResponseHeader<Self, M>where
Self: Sized,
fn append_response_header<M>(
self,
header_name: HeaderName,
make: M,
) -> SetResponseHeader<Self, M>where
Self: Sized,
§fn insert_response_header_if_not_present<M>(
self,
header_name: HeaderName,
make: M,
) -> SetResponseHeader<Self, M>where
Self: Sized,
fn insert_response_header_if_not_present<M>(
self,
header_name: HeaderName,
make: M,
) -> SetResponseHeader<Self, M>where
Self: Sized,
§fn catch_panic(self) -> CatchPanic<Self, DefaultResponseForPanic>where
Self: Sized,
fn catch_panic(self) -> CatchPanic<Self, DefaultResponseForPanic>where
Self: Sized,
500 Internal Server responses. Read moreSource§impl<T> Tried for T
impl<T> Tried for T
Source§fn try_add(self, rhs: Self) -> Result<Self, Error>where
Self: Sized + CheckedAdd,
fn try_add(self, rhs: Self) -> Result<Self, Error>where
Self: Sized + CheckedAdd,
rhs with checked arithmetic. Read moreSource§fn try_sub(self, rhs: Self) -> Result<Self, Error>where
Self: Sized + CheckedSub,
fn try_sub(self, rhs: Self) -> Result<Self, Error>where
Self: Sized + CheckedSub,
rhs with checked arithmetic. Read moreSource§fn try_mul(self, rhs: Self) -> Result<Self, Error>where
Self: Sized + CheckedMul,
fn try_mul(self, rhs: Self) -> Result<Self, Error>where
Self: Sized + CheckedMul,
rhs with checked arithmetic. Read more