chore: format logs
This commit is contained in:
parent
507d19628b
commit
eb4814738b
|
|
@ -270,22 +270,16 @@ where
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let mut merged_msg = vec![];
|
||||||
// If the message can merge other messages, try to merge the next message until the
|
// If the message can merge other messages, try to merge the next message until the
|
||||||
// message is not mergeable.
|
// message is not mergeable.
|
||||||
if sending_msg.can_merge() {
|
if sending_msg.can_merge() {
|
||||||
let mut merged_msg = vec![];
|
|
||||||
while let Some(pending_msg) = pending_msg_queue.pop() {
|
while let Some(pending_msg) = pending_msg_queue.pop() {
|
||||||
// If the message is not mergeable, push the message back to the queue and break the loop.
|
// If the message is not mergeable, push the message back to the queue and break the loop.
|
||||||
match sending_msg.merge(&pending_msg, &self.config.maximum_payload_size) {
|
match sending_msg.merge(&pending_msg, &self.config.maximum_payload_size) {
|
||||||
Ok(continue_merge) => {
|
Ok(continue_merge) => {
|
||||||
merged_msg.push(pending_msg.msg_id());
|
merged_msg.push(pending_msg.msg_id());
|
||||||
if !continue_merge {
|
if !continue_merge {
|
||||||
event!(
|
|
||||||
tracing::Level::TRACE,
|
|
||||||
"merge: {:?}, len: {}",
|
|
||||||
merged_msg,
|
|
||||||
sending_msg.get_msg().length()
|
|
||||||
);
|
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|
@ -300,9 +294,19 @@ where
|
||||||
|
|
||||||
sending_msg.set_ret(tx);
|
sending_msg.set_ret(tx);
|
||||||
sending_msg.set_state(self.uid, MessageState::Processing);
|
sending_msg.set_state(self.uid, MessageState::Processing);
|
||||||
|
|
||||||
let _ = self.state_notifier.send(SinkState::Syncing);
|
let _ = self.state_notifier.send(SinkState::Syncing);
|
||||||
let collab_msg = sending_msg.get_msg().clone();
|
let collab_msg = sending_msg.get_msg().clone();
|
||||||
pending_msg_queue.push(sending_msg);
|
pending_msg_queue.push(sending_msg);
|
||||||
|
|
||||||
|
if !merged_msg.is_empty() {
|
||||||
|
event!(
|
||||||
|
tracing::Level::DEBUG,
|
||||||
|
"merge: {:?}, len: {}",
|
||||||
|
merged_msg,
|
||||||
|
collab_msg.length()
|
||||||
|
);
|
||||||
|
}
|
||||||
collab_msg
|
collab_msg
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
@ -456,7 +460,7 @@ impl Default for SinkConfig {
|
||||||
fn default() -> Self {
|
fn default() -> Self {
|
||||||
Self {
|
Self {
|
||||||
send_timeout: Duration::from_secs(DEFAULT_SYNC_TIMEOUT),
|
send_timeout: Duration::from_secs(DEFAULT_SYNC_TIMEOUT),
|
||||||
maximum_payload_size: 4096,
|
maximum_payload_size: 1024 * 64,
|
||||||
strategy: SinkStrategy::ASAP,
|
strategy: SinkStrategy::ASAP,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue