From fd0f45897874b7d77e17cbcab86a31d2d71146b2 Mon Sep 17 00:00:00 2001 From: timedout Date: Sun, 7 Jun 2026 02:05:37 +0100 Subject: [PATCH] style: Tidy up --- src/service/rooms/event_handler/fetch_prev.rs | 5 ++- .../event_handler/handle_incoming_pdu.rs | 36 ++++++++++--------- 2 files changed, 21 insertions(+), 20 deletions(-) diff --git a/src/service/rooms/event_handler/fetch_prev.rs b/src/service/rooms/event_handler/fetch_prev.rs index e49f6e7d9..e9a883c33 100644 --- a/src/service/rooms/event_handler/fetch_prev.rs +++ b/src/service/rooms/event_handler/fetch_prev.rs @@ -1,7 +1,7 @@ use std::{collections::HashMap, time::Instant}; use conduwuit::{ - Event, PduEvent, debug, debug_info, debug_warn, info, trace, + Event, PduEvent, debug, debug_info, debug_warn, trace, utils::{BoolExt, IterStream, stream::BroadbandExt}, }; use futures::StreamExt; @@ -84,13 +84,12 @@ pub(super) async fn fetch_prevs( let job_start = Instant::now(); trace!("Starting to persist {} prev events", to_persist.len()); for (i, event_id) in to_persist.iter().enumerate() { - info!( + debug!( elapsed=?start.elapsed(), "[TODO] Persisting fetched prev event: {event_id} ({}/{})", i.saturating_add(1), to_persist.len(), ); - debug_info!(elapsed=?start.elapsed(), "Persisting fetched prev event {event_id}"); let obj = mapped.get(event_id).cloned().unwrap(); let persist_start = Instant::now(); match self diff --git a/src/service/rooms/event_handler/handle_incoming_pdu.rs b/src/service/rooms/event_handler/handle_incoming_pdu.rs index 83297d0a3..8037fd8da 100644 --- a/src/service/rooms/event_handler/handle_incoming_pdu.rs +++ b/src/service/rooms/event_handler/handle_incoming_pdu.rs @@ -1,7 +1,8 @@ use std::{collections::BTreeMap, time::Instant}; use conduwuit::{ - Err, Event, PduEvent, Result, debug_error, debug_info, defer, err, error, info, trace, warn, + Err, Event, PduEvent, Result, debug, debug_error, debug_info, defer, err, error, info, + result::DebugInspect, trace, warn, }; use futures::{ FutureExt, @@ -226,38 +227,39 @@ pub async fn handle_incoming_pdu<'a>( .remove(room_id); }} - info!("[TODO] Handling PDU as outlier"); let (incoming_pdu, val) = self .handle_outlier_pdu(origin, create_event, event_id, room_id, value, false) - .await - .inspect_err(|e| error!("[TODO] Failed to handle outlier PDU: {e:?}"))?; - info!("[TODO] Finished handling PDU as outlier"); + .await?; + // 8. if not timeline event: stop if !is_timeline_event { - info!("[TODO] Not upgrading PDU"); return Ok(None); } - // Skip old events - // let first_ts_in_room = self - // .services - // .timeline - // .first_pdu_in_room(room_id) - // .await? - // .origin_server_ts(); + // Skip events sent before we joined (they need to be persisted as backfilled + // events, not timeline events, which is handled elsewhere). + let first_ts_in_room = self + .services + .timeline + .first_pdu_in_room(room_id) + .await? + .origin_server_ts(); + if incoming_pdu.origin_server_ts() < first_ts_in_room { + return Ok(None); + } // 9. Fetch any missing prev events doing all checks listed here starting at 1. // These are timeline events - info!("[TODO] Fetching prev events"); + debug!("Fetching and persisting any missing prev events"); self.fetch_prevs(room_id, create_event, &incoming_pdu, origin) .await - .inspect_err(|e| error!("[TODO] Failed to fetch_prevs: {e:?}"))?; + .debug_inspect_err(|e| { + error!("Failed to fetch and persist incoming event's prev_events: {e:?}"); + })?; - info!("[TODO] Finished fetching prev events, attempting to upgrade"); // Done with prev events, now handling the incoming event self.upgrade_outlier_to_timeline_pdu(incoming_pdu, val, create_event, origin, room_id) .await - .inspect_err(|e| error!("[TODO] Failed to upgrade outlier to timeline pdu: {e:?}")) } }