Misc debug and trace log tweaks.
Signed-off-by: Jason Volk <jason@zemos.net>
This commit is contained in:
@@ -27,6 +27,11 @@ use tuwunel_core::{
|
||||
/// c. Ask origin server over federation
|
||||
/// d. TODO: Ask other servers over federation?
|
||||
#[implement(super::Service)]
|
||||
#[tracing::instrument(
|
||||
level = "debug",
|
||||
skip_all,
|
||||
fields(%origin),
|
||||
)]
|
||||
pub(super) async fn fetch_auth<'a, Events>(
|
||||
&self,
|
||||
origin: &ServerName,
|
||||
@@ -96,6 +101,12 @@ where
|
||||
}
|
||||
|
||||
#[implement(super::Service)]
|
||||
#[tracing::instrument(
|
||||
name = "chain",
|
||||
level = "trace",
|
||||
skip_all,
|
||||
fields(%event_id),
|
||||
)]
|
||||
async fn fetch_auth_chain(
|
||||
&self,
|
||||
origin: &ServerName,
|
||||
@@ -127,7 +138,7 @@ async fn fetch_auth_chain(
|
||||
start: Duration::from_secs(2 * 60),
|
||||
end: Duration::from_secs(60 * 60 * 8),
|
||||
}) {
|
||||
debug_warn!("Backing off from {next_id}");
|
||||
debug_warn!("Backed off from {next_id}");
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -144,6 +155,7 @@ async fn fetch_auth_chain(
|
||||
.await
|
||||
.inspect_err(|e| debug_error!("Failed to fetch event {next_id}: {e}"))
|
||||
else {
|
||||
debug_warn!("Backing off from {next_id}");
|
||||
self.back_off(&next_id);
|
||||
continue;
|
||||
};
|
||||
|
||||
@@ -18,7 +18,7 @@ use super::check_room_id;
|
||||
|
||||
#[implement(super::Service)]
|
||||
#[tracing::instrument(
|
||||
level = "debug",
|
||||
level = "debug",
|
||||
skip_all,
|
||||
fields(%origin),
|
||||
)]
|
||||
|
||||
@@ -45,7 +45,7 @@ use crate::rooms::timeline::RawPduId;
|
||||
level = INFO_SPAN_LEVEL,
|
||||
skip_all,
|
||||
fields(%room_id, %event_id),
|
||||
ret(Debug),
|
||||
ret(level = "debug"),
|
||||
)]
|
||||
pub async fn handle_incoming_pdu<'a>(
|
||||
&'a self,
|
||||
@@ -57,7 +57,7 @@ pub async fn handle_incoming_pdu<'a>(
|
||||
) -> Result<Option<(RawPduId, bool)>> {
|
||||
// 1. Skip the PDU if we already have it as a timeline event
|
||||
if let Ok(pdu_id) = self.services.timeline.get_pdu_id(event_id).await {
|
||||
trace!(?event_id, "exists");
|
||||
debug!(?pdu_id, "Exists.");
|
||||
return Ok(Some((pdu_id, false)));
|
||||
}
|
||||
|
||||
@@ -116,6 +116,10 @@ pub async fn handle_incoming_pdu<'a>(
|
||||
|
||||
// 8. if not timeline event: stop
|
||||
if !is_timeline_event {
|
||||
debug!(
|
||||
kind = ?incoming_pdu.event_type(),
|
||||
"Not a timeline event.",
|
||||
);
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
@@ -128,6 +132,11 @@ pub async fn handle_incoming_pdu<'a>(
|
||||
.origin_server_ts();
|
||||
|
||||
if incoming_pdu.origin_server_ts() < first_ts_in_room {
|
||||
debug!(
|
||||
origin_server_ts = ?incoming_pdu.origin_server_ts(),
|
||||
?first_ts_in_room,
|
||||
"Skipping old event."
|
||||
);
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
@@ -137,11 +146,11 @@ pub async fn handle_incoming_pdu<'a>(
|
||||
.fetch_prev(origin, room_id, incoming_pdu.prev_events(), &room_version, first_ts_in_room)
|
||||
.await?;
|
||||
|
||||
debug!(
|
||||
events = ?sorted_prev_events,
|
||||
trace!(
|
||||
events = sorted_prev_events.len(),
|
||||
event_ids = ?sorted_prev_events,
|
||||
"Handling previous events"
|
||||
);
|
||||
|
||||
sorted_prev_events
|
||||
.iter()
|
||||
.try_stream()
|
||||
|
||||
@@ -23,6 +23,8 @@ pub(super) async fn handle_outlier_pdu(
|
||||
room_version: &RoomVersionId,
|
||||
auth_events_known: bool,
|
||||
) -> Result<(PduEvent, CanonicalJsonObject)> {
|
||||
debug!(?event_id, ?auth_events_known, "handle outlier");
|
||||
|
||||
// 1. Remove unsigned field
|
||||
pdu_json.remove("unsigned");
|
||||
|
||||
|
||||
@@ -6,7 +6,7 @@ use std::{
|
||||
use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt, future::try_join};
|
||||
use ruma::{OwnedEventId, RoomId, RoomVersionId};
|
||||
use tuwunel_core::{
|
||||
Result, apply, err, implement,
|
||||
Result, apply, debug, debug_warn, err, implement,
|
||||
matrix::{Event, StateMap, state_res::AuthSet},
|
||||
ref_at, trace,
|
||||
utils::{
|
||||
@@ -42,11 +42,14 @@ where
|
||||
.services
|
||||
.state
|
||||
.pdu_shortstatehash(prev_event_id)
|
||||
.inspect_err(|e| debug_warn!(?prev_event_id, "Missing state at prev_event: {e}"))
|
||||
.await
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
debug!(?prev_event_id, ?prev_event_sstatehash, "Resolving state at prev_event.");
|
||||
|
||||
let prev_event = self
|
||||
.services
|
||||
.timeline
|
||||
@@ -62,6 +65,13 @@ where
|
||||
|
||||
let (prev_event, mut state) = try_join(prev_event, state).await?;
|
||||
|
||||
debug!(
|
||||
?prev_event_id,
|
||||
?prev_event_sstatehash,
|
||||
state_ids = state.len(),
|
||||
"Resolved state at prev_event.",
|
||||
);
|
||||
|
||||
if let Some(state_key) = prev_event.state_key() {
|
||||
let prev_event_type = prev_event.event_type().to_cow_str().into();
|
||||
|
||||
@@ -73,6 +83,14 @@ where
|
||||
|
||||
state.insert(shortstatekey, prev_event.event_id().into());
|
||||
// Now it's the state after the pdu
|
||||
debug!(
|
||||
?prev_event_id,
|
||||
?prev_event_type,
|
||||
?prev_event_sstatehash,
|
||||
?shortstatekey,
|
||||
state_ids = state.len(),
|
||||
"Added prev_event to state.",
|
||||
);
|
||||
}
|
||||
|
||||
debug_assert!(!state.is_empty(), "should be returning None for empty HashMap result");
|
||||
@@ -101,14 +119,16 @@ where
|
||||
.prev_events()
|
||||
.try_stream()
|
||||
.broad_and_then(|prev_event_id| {
|
||||
let prev_event = self.services.timeline.get_pdu(prev_event_id);
|
||||
|
||||
let sstatehash = self
|
||||
.services
|
||||
.state
|
||||
.pdu_shortstatehash(prev_event_id);
|
||||
|
||||
try_join(sstatehash, prev_event)
|
||||
let prev_event = self.services.timeline.get_pdu(prev_event_id);
|
||||
|
||||
try_join(sstatehash, prev_event).inspect_err(move |e| {
|
||||
debug_warn!(?prev_event_id, "Missing state at prev_event: {e}");
|
||||
})
|
||||
})
|
||||
.try_collect::<HashMap<_, _>>()
|
||||
.await
|
||||
@@ -133,12 +153,12 @@ where
|
||||
trace!("Resolving state");
|
||||
let Ok(new_state) = self
|
||||
.state_resolution(room_version_id, fork_states, auth_chain_sets)
|
||||
.inspect_ok(|_| trace!("State resolution done."))
|
||||
.await
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
trace!("State resolution done.");
|
||||
new_state
|
||||
.into_iter()
|
||||
.stream()
|
||||
@@ -149,7 +169,8 @@ where
|
||||
.map(move |shortstatekey| (shortstatekey, event_id))
|
||||
.await
|
||||
})
|
||||
.collect()
|
||||
.collect::<HashMap<_, _>>()
|
||||
.inspect(|state| trace!(state = state.len(), "Created shortstatekeys."))
|
||||
.map(Some)
|
||||
.map(Ok)
|
||||
.await
|
||||
@@ -192,6 +213,13 @@ where
|
||||
.collect()
|
||||
.await;
|
||||
|
||||
trace!(
|
||||
prev_event = ?prev_event.event_id(),
|
||||
?sstatehash,
|
||||
leaf_states = leaf_state_after_event.len(),
|
||||
"leaf state after event"
|
||||
);
|
||||
|
||||
let starting_events = leaf_state_after_event
|
||||
.iter()
|
||||
.map(ref_at!(1))
|
||||
|
||||
@@ -19,7 +19,12 @@ use crate::rooms::{
|
||||
};
|
||||
|
||||
#[implement(super::Service)]
|
||||
#[tracing::instrument(name = "upgrade", level = "debug", skip_all, ret(Debug))]
|
||||
#[tracing::instrument(
|
||||
name = "upgrade",
|
||||
level = "debug",
|
||||
skip_all,
|
||||
ret(level = "debug")
|
||||
)]
|
||||
pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
&self,
|
||||
origin: &ServerName,
|
||||
@@ -30,13 +35,14 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
create_event_id: &EventId,
|
||||
) -> Result<Option<(RawPduId, bool)>> {
|
||||
// Skip the PDU if we already have it as a timeline event
|
||||
if let Ok(pduid) = self
|
||||
if let Ok(pdu_id) = self
|
||||
.services
|
||||
.timeline
|
||||
.get_pdu_id(incoming_pdu.event_id())
|
||||
.await
|
||||
{
|
||||
return Ok(Some((pduid, false)));
|
||||
debug!(?pdu_id, "Exists.");
|
||||
return Ok(Some((pdu_id, false)));
|
||||
}
|
||||
|
||||
if self
|
||||
@@ -48,17 +54,19 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
return Err!(Request(InvalidParam("Event has been soft failed")));
|
||||
}
|
||||
|
||||
debug!("Upgrading to timeline pdu");
|
||||
trace!("Upgrading to timeline pdu");
|
||||
|
||||
let timer = Instant::now();
|
||||
let room_rules = room_version::rules(room_version)?;
|
||||
|
||||
trace!(format = ?room_rules.event_format, "Checking format");
|
||||
state_res::check_pdu_format(&val, &room_rules.event_format)?;
|
||||
|
||||
// 10. Fetch missing state and auth chain events by calling /state_ids at
|
||||
// backwards extremities doing all the checks in this list starting at 1.
|
||||
// These are not timeline events.
|
||||
trace!("Resolving state at event");
|
||||
|
||||
debug!("Resolving state at event");
|
||||
let mut state_at_incoming_event = if incoming_pdu.prev_events().count() == 1 {
|
||||
self.state_at_incoming_degree_one(&incoming_pdu)
|
||||
.await?
|
||||
@@ -78,8 +86,8 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
let state_at_incoming_event =
|
||||
state_at_incoming_event.expect("we always set this to some above");
|
||||
|
||||
debug!("Performing auth check");
|
||||
// 11. Check the auth of the event passes based on the state of the event
|
||||
|
||||
let state_fetch = async |k: StateEventType, s: StateKey| {
|
||||
let shortstatekey = self
|
||||
.services
|
||||
@@ -99,9 +107,11 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
};
|
||||
|
||||
let event_fetch = async |event_id: OwnedEventId| self.event_fetch(&event_id).await;
|
||||
|
||||
trace!("Performing auth check");
|
||||
state_res::auth_check(&room_rules, &incoming_pdu, &event_fetch, &state_fetch).await?;
|
||||
|
||||
debug!("Gathering auth events");
|
||||
trace!("Gathering auth events");
|
||||
let auth_events = self
|
||||
.services
|
||||
.state
|
||||
@@ -123,10 +133,11 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
.ok_or_else(|| err!(Request(NotFound("state event not found"))))
|
||||
};
|
||||
|
||||
trace!("Performing auth check");
|
||||
state_res::auth_check(&room_rules, &incoming_pdu, &event_fetch, &state_fetch).await?;
|
||||
|
||||
// Soft fail check before doing state res
|
||||
debug!("Performing soft-fail check");
|
||||
trace!("Performing soft-fail check");
|
||||
let soft_fail = match incoming_pdu.redacts_id(room_version) {
|
||||
| None => false,
|
||||
| Some(redact_id) =>
|
||||
@@ -138,7 +149,6 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
};
|
||||
|
||||
// 13. Use state resolution to find new room state
|
||||
|
||||
// We start looking at current room state now, so lets lock the room
|
||||
trace!("Locking the room");
|
||||
let state_lock = self.services.state.mutex.lock(room_id).await;
|
||||
@@ -170,11 +180,12 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
.await;
|
||||
|
||||
debug!(
|
||||
"Retained {} extremities checked against {} prev_events",
|
||||
extremities.len(),
|
||||
incoming_pdu.prev_events().count()
|
||||
retained = extremities.len(),
|
||||
prev_events = incoming_pdu.prev_events().count(),
|
||||
"Retained extremities checked against prev_events.",
|
||||
);
|
||||
|
||||
trace!("Compressing state...");
|
||||
let state_ids_compressed: Arc<CompressedState> = self
|
||||
.services
|
||||
.state_compressor
|
||||
@@ -188,34 +199,49 @@ pub(super) async fn upgrade_outlier_to_timeline_pdu(
|
||||
.await;
|
||||
|
||||
if incoming_pdu.state_key().is_some() {
|
||||
debug!("Event is a state-event. Deriving new room state");
|
||||
|
||||
// We also add state after incoming event to the fork states
|
||||
let mut state_after = state_at_incoming_event.clone();
|
||||
if let Some(state_key) = incoming_pdu.state_key() {
|
||||
let event_id = incoming_pdu.event_id();
|
||||
let event_type = incoming_pdu.kind();
|
||||
let shortstatekey = self
|
||||
.services
|
||||
.short
|
||||
.get_or_create_shortstatekey(&incoming_pdu.kind().to_string().into(), state_key)
|
||||
.get_or_create_shortstatekey(&event_type.to_string().into(), state_key)
|
||||
.await;
|
||||
|
||||
let event_id = incoming_pdu.event_id();
|
||||
state_after.insert(shortstatekey, event_id.to_owned());
|
||||
// Now it's the state after the event.
|
||||
debug!(
|
||||
?event_id,
|
||||
?event_type,
|
||||
?state_key,
|
||||
?shortstatekey,
|
||||
state_after = state_after.len(),
|
||||
"Adding event to state."
|
||||
);
|
||||
}
|
||||
|
||||
trace!("Resolving new room state.");
|
||||
let new_room_state = self
|
||||
.resolve_state(room_id, room_version, state_after)
|
||||
.boxed()
|
||||
.await?;
|
||||
|
||||
// Set the new room state to the resolved state
|
||||
debug!("Forcing new room state");
|
||||
trace!("Saving resolved state.");
|
||||
let HashSetCompressStateEvent { shortstatehash, added, removed } = self
|
||||
.services
|
||||
.state_compressor
|
||||
.save_state(room_id, new_room_state)
|
||||
.await?;
|
||||
|
||||
debug!(
|
||||
?shortstatehash,
|
||||
added = added.len(),
|
||||
removed = removed.len(),
|
||||
"Forcing new room state."
|
||||
);
|
||||
self.services
|
||||
.state
|
||||
.force_state(room_id, shortstatehash, added, removed, &state_lock)
|
||||
|
||||
Reference in New Issue
Block a user