use std::sync::Arc;
use std::time::Instant;
use bytes::Buf;
use bytes::BufMut;
use cheetah_string::CheetahString;
use dashmap::DashMap;
use rocketmq_common::common::config::TopicConfig;
use rocketmq_common::common::message::message_batch::MessageExtBatch;
use rocketmq_common::common::message::message_ext_broker_inner::MessageExtBrokerInner;
use rocketmq_common::common::sys_flag::message_sys_flag::MessageSysFlag;
use rocketmq_common::utils::message_utils;
use rocketmq_common::CRC32Utils::crc32;
use rocketmq_common::MessageDecoder::create_crc32;
use rocketmq_common::MessageUtils::build_batch_message_id;
use rocketmq_rust::ArcMut;
use rocketmq_rust::SyncUnsafeCellWrapper;
use tracing::error;
use crate::base::message_result::AppendMessageResult;
use crate::base::message_status_enum::AppendMessageStatus;
use crate::base::put_message_context::PutMessageContext;
use crate::config::message_store_config::MessageStoreConfig;
use crate::log_file::commit_log::get_message_num;
use crate::log_file::commit_log::CommitLog;
use crate::log_file::commit_log::BLANK_MAGIC_CODE;
use crate::log_file::commit_log::CRC32_RESERVED_LEN;
use crate::log_file::mapped_file::MappedFile;
pub trait AppendMessageCallback {
fn do_append<MF: MappedFile>(
&self,
file_from_offset: i64,
mapped_file: &MF,
max_blank: i32,
msg: &mut MessageExtBrokerInner,
put_message_context: &PutMessageContext,
) -> AppendMessageResult;
fn do_append_batch<MF: MappedFile>(
&self,
file_from_offset: i64,
mapped_file: &MF,
max_blank: i32,
msg: &mut MessageExtBatch,
put_message_context: &mut PutMessageContext,
enabled_append_prop_crc: bool,
) -> AppendMessageResult;
fn do_append_zerocopy<MF: MappedFile>(
&self,
file_from_offset: i64,
mapped_file: &MF,
max_blank: i32,
msg: &mut MessageExtBrokerInner,
put_message_context: &PutMessageContext,
) -> AppendMessageResult {
self.do_append(file_from_offset, mapped_file, max_blank, msg, put_message_context)
}
}
const END_FILE_MIN_BLANK_LENGTH: i32 = 4 + 4;
pub struct DefaultAppendMessageCallback {
msg_store_item_memory: SyncUnsafeCellWrapper<bytes::BytesMut>,
crc32_reserved_length: i32,
message_store_config: Arc<MessageStoreConfig>,
topic_config_table: Arc<DashMap<CheetahString, ArcMut<TopicConfig>>>,
}
impl DefaultAppendMessageCallback {
pub fn new(
message_store_config: Arc<MessageStoreConfig>,
topic_config_table: Arc<DashMap<CheetahString, ArcMut<TopicConfig>>>,
) -> Self {
let crc32_reserved_length = if message_store_config.enabled_append_prop_crc {
CRC32_RESERVED_LEN
} else {
0
};
Self {
msg_store_item_memory: SyncUnsafeCellWrapper::new(bytes::BytesMut::with_capacity(
END_FILE_MIN_BLANK_LENGTH as usize,
)),
crc32_reserved_length,
message_store_config,
topic_config_table,
}
}
}
impl AppendMessageCallback for DefaultAppendMessageCallback {
fn do_append<MF: MappedFile>(
&self,
file_from_offset: i64,
mapped_file: &MF,
max_blank: i32,
msg_inner: &mut MessageExtBrokerInner,
put_message_context: &PutMessageContext,
) -> AppendMessageResult {
let mut pre_encode_buffer = msg_inner.encoded_buff.take().unwrap(); let is_multi_dispatch_msg =
self.message_store_config.enable_multi_dispatch && CommitLog::is_multi_dispatch_msg(msg_inner);
if is_multi_dispatch_msg {
unimplemented!("Multi dispatch message is not supported yet");
}
let msg_len = i32::from_be_bytes(pre_encode_buffer[0..4].try_into().unwrap());
let wrote_offset = file_from_offset + mapped_file.get_wrote_position() as i64;
let addr = msg_inner.message_ext_inner.store_host;
let msg_id_supplier = move || -> String { message_utils::build_message_id(addr, wrote_offset) };
let mut queue_offset = msg_inner.queue_offset();
let message_num = get_message_num(&self.topic_config_table, msg_inner);
if let MessageSysFlag::TRANSACTION_PREPARED_TYPE | MessageSysFlag::TRANSACTION_ROLLBACK_TYPE =
MessageSysFlag::get_transaction_value(msg_inner.sys_flag())
{
queue_offset = 0;
}
if (msg_len + END_FILE_MIN_BLANK_LENGTH) > max_blank {
let bytes = self.msg_store_item_memory.mut_from_ref();
bytes.clear();
bytes.put_i32(max_blank);
bytes.put_i32(BLANK_MAGIC_CODE);
let instant = Instant::now();
mapped_file.write_bytes_segment(bytes.as_ref(), wrote_offset as usize, 0, bytes.len());
return AppendMessageResult {
status: AppendMessageStatus::EndOfFile,
wrote_offset,
wrote_bytes: max_blank,
store_timestamp: msg_inner.store_timestamp(),
logics_offset: queue_offset,
msg_num: message_num as i32,
msg_id_supplier: Some(Arc::new(msg_id_supplier)),
page_cache_rt: instant.elapsed().as_millis() as i64,
..Default::default()
};
}
let mut pos = 20; pre_encode_buffer[pos..(pos + 8)].copy_from_slice(&queue_offset.to_be_bytes()); pos += 8;
pre_encode_buffer[pos..(pos + 8)].copy_from_slice(&wrote_offset.to_be_bytes()); let ip_len = if msg_inner.sys_flag() & MessageSysFlag::BORNHOST_V6_FLAG == 0 {
4 + 4
} else {
16 + 4
};
pos += 8 + 4 + 8 + ip_len;
pre_encode_buffer[pos..(pos + 8)].copy_from_slice(&msg_inner.store_timestamp().to_be_bytes());
if self.message_store_config.enabled_append_prop_crc {
let check_size = msg_len - self.crc32_reserved_length;
let crc32 = crc32(&pre_encode_buffer[..check_size as usize]);
create_crc32(&mut pre_encode_buffer[check_size as usize..msg_len as usize], crc32);
}
let instant = Instant::now();
mapped_file.append_message_bytes_no_position_update_ref(pre_encode_buffer.chunk());
AppendMessageResult {
status: AppendMessageStatus::PutOk,
wrote_offset,
wrote_bytes: msg_len,
store_timestamp: msg_inner.store_timestamp(),
logics_offset: queue_offset,
msg_num: message_num as i32,
msg_id_supplier: Some(Arc::new(Box::new(msg_id_supplier))),
page_cache_rt: instant.elapsed().as_millis() as i64,
..Default::default()
}
}
fn do_append_batch<MF: MappedFile>(
&self,
file_from_offset: i64,
mapped_file: &MF,
max_blank: i32,
msg_batch: &mut MessageExtBatch,
put_message_context: &mut PutMessageContext,
enabled_append_prop_crc: bool,
) -> AppendMessageResult {
let wrote_offset = file_from_offset + mapped_file.get_wrote_position() as i64;
let queue_offset = msg_batch.message_ext_broker_inner.queue_offset();
let begin_queue_offset = queue_offset;
let begin_time_mills = Instant::now();
let mut messages_byte_buffer = msg_batch.encoded_buff.take().unwrap();
let sys_flag = msg_batch.message_ext_broker_inner.sys_flag();
let born_host_length = if sys_flag & MessageSysFlag::BORNHOST_V6_FLAG == 0 {
4 + 4
} else {
16 + 4
};
let store_host_length = if sys_flag & MessageSysFlag::STOREHOSTADDRESS_V6_FLAG == 0 {
4 + 4
} else {
16 + 4
};
let addr = msg_batch.message_ext_broker_inner.store_host();
let batch_size = put_message_context.get_batch_size();
let phy_ops = put_message_context.get_phy_pos().to_vec();
let msg_id_supplier =
move || -> String { build_batch_message_id(addr, store_host_length, batch_size as usize, &phy_ops) };
let mut total_msg_len = 0;
let mut msg_num = 0;
let mut msg_pos = 0;
let mut index = 0;
while total_msg_len < messages_byte_buffer.len() as i32 {
let msg_len = i32::from_be_bytes(
messages_byte_buffer[total_msg_len as usize..(total_msg_len + 4) as usize]
.try_into()
.unwrap(),
);
total_msg_len += msg_len;
if total_msg_len + END_FILE_MIN_BLANK_LENGTH > max_blank {
let bytes = self.msg_store_item_memory.mut_from_ref();
bytes.clear();
bytes.put_i32(max_blank);
bytes.put_i32(BLANK_MAGIC_CODE);
mapped_file.write_bytes_segment(bytes.as_ref(), wrote_offset as usize, 0, bytes.len());
return AppendMessageResult {
status: AppendMessageStatus::EndOfFile,
wrote_offset,
wrote_bytes: max_blank,
msg_id_supplier: Some(Arc::new(Box::new(msg_id_supplier))),
store_timestamp: msg_batch.message_ext_broker_inner.store_timestamp(),
logics_offset: begin_queue_offset,
page_cache_rt: begin_time_mills.elapsed().as_millis() as i64,
..Default::default()
};
}
let mut pos = msg_pos + 20;
messages_byte_buffer[pos..(pos + 8)].copy_from_slice(&queue_offset.to_be_bytes());
pos += 8;
let phy_pos = wrote_offset + total_msg_len as i64 - msg_len as i64;
messages_byte_buffer[pos..(pos + 8)].copy_from_slice(&phy_pos.to_be_bytes());
pos += 8 + 4 + 8 + born_host_length;
messages_byte_buffer[pos..(pos + 8)]
.copy_from_slice(&msg_batch.message_ext_broker_inner.store_timestamp().to_be_bytes());
if enabled_append_prop_crc {
let _check_size = msg_len - self.crc32_reserved_length;
}
put_message_context.get_phy_pos_mut()[index] = phy_pos;
msg_num += 1;
msg_pos += msg_len as usize;
index += 1;
}
let bytes = messages_byte_buffer.freeze();
mapped_file.append_message_bytes_no_position_update(&bytes);
AppendMessageResult {
status: AppendMessageStatus::PutOk,
wrote_offset,
wrote_bytes: total_msg_len,
msg_id_supplier: Some(Arc::new(Box::new(msg_id_supplier))),
store_timestamp: msg_batch.message_ext_broker_inner.store_timestamp(),
logics_offset: begin_queue_offset,
page_cache_rt: begin_time_mills.elapsed().as_millis() as i64,
msg_num,
..Default::default()
}
}
fn do_append_zerocopy<MF: MappedFile>(
&self,
file_from_offset: i64,
mapped_file: &MF,
max_blank: i32,
msg_inner: &mut MessageExtBrokerInner,
put_message_context: &PutMessageContext,
) -> AppendMessageResult {
let pre_encode_buffer = msg_inner.encoded_buff.take().unwrap();
let is_multi_dispatch_msg =
self.message_store_config.enable_multi_dispatch && CommitLog::is_multi_dispatch_msg(msg_inner);
if is_multi_dispatch_msg {
msg_inner.encoded_buff = Some(pre_encode_buffer);
return self.do_append(file_from_offset, mapped_file, max_blank, msg_inner, put_message_context);
}
let msg_len = i32::from_be_bytes(pre_encode_buffer[0..4].try_into().unwrap());
let wrote_offset = file_from_offset + mapped_file.get_wrote_position() as i64;
let addr = msg_inner.message_ext_inner.store_host;
let msg_id_supplier = move || -> String { message_utils::build_message_id(addr, wrote_offset) };
let mut queue_offset = msg_inner.queue_offset();
let message_num = get_message_num(&self.topic_config_table, msg_inner);
if let MessageSysFlag::TRANSACTION_PREPARED_TYPE | MessageSysFlag::TRANSACTION_ROLLBACK_TYPE =
MessageSysFlag::get_transaction_value(msg_inner.sys_flag())
{
queue_offset = 0;
}
if (msg_len + END_FILE_MIN_BLANK_LENGTH) > max_blank {
let bytes = self.msg_store_item_memory.mut_from_ref();
bytes.clear();
bytes.put_i32(max_blank);
bytes.put_i32(BLANK_MAGIC_CODE);
let instant = Instant::now();
mapped_file.write_bytes_segment(bytes.as_ref(), wrote_offset as usize, 0, bytes.len());
return AppendMessageResult {
status: AppendMessageStatus::EndOfFile,
wrote_offset,
wrote_bytes: max_blank,
store_timestamp: msg_inner.store_timestamp(),
logics_offset: queue_offset,
msg_num: message_num as i32,
msg_id_supplier: Some(Arc::new(msg_id_supplier)),
page_cache_rt: instant.elapsed().as_millis() as i64,
..Default::default()
};
}
let instant = Instant::now();
if let Some((buffer, _pos)) = mapped_file.get_direct_write_buffer(msg_len as usize) {
buffer[..msg_len as usize].copy_from_slice(&pre_encode_buffer[..msg_len as usize]);
let mut pos = 20;
buffer[pos..pos + 8].copy_from_slice(&queue_offset.to_be_bytes());
pos += 8;
buffer[pos..pos + 8].copy_from_slice(&wrote_offset.to_be_bytes());
let ip_len = if msg_inner.sys_flag() & MessageSysFlag::BORNHOST_V6_FLAG == 0 {
4 + 4
} else {
16 + 4
};
pos += 8 + 4 + 8 + ip_len;
buffer[pos..pos + 8].copy_from_slice(&msg_inner.store_timestamp().to_be_bytes());
if self.message_store_config.enabled_append_prop_crc {
let check_size = msg_len - self.crc32_reserved_length;
let crc32 = crc32(&buffer[..check_size as usize]);
create_crc32(&mut buffer[check_size as usize..msg_len as usize], crc32);
}
if !mapped_file.commit_direct_write(msg_len as usize) {
error!("Failed to commit zero-copy write");
return AppendMessageResult {
status: AppendMessageStatus::UnknownError,
..Default::default()
};
}
AppendMessageResult {
status: AppendMessageStatus::PutOk,
wrote_offset,
wrote_bytes: msg_len,
store_timestamp: msg_inner.store_timestamp(),
logics_offset: queue_offset,
msg_num: message_num as i32,
msg_id_supplier: Some(Arc::new(Box::new(msg_id_supplier))),
page_cache_rt: instant.elapsed().as_millis() as i64,
..Default::default()
}
} else {
msg_inner.encoded_buff = Some(pre_encode_buffer);
self.do_append(file_from_offset, mapped_file, max_blank, msg_inner, put_message_context)
}
}
}