#![allow(unused_variables)]
use std::any::Any;
use std::collections::HashMap;
use std::fs;
use std::path::Path;
use std::str::FromStr;
use std::sync::Arc;
use bytes::BufMut;
use bytes::Bytes;
use bytes::BytesMut;
use cheetah_string::CheetahString;
use futures_util::future::join_all;
use rocketmq_common::common::attribute::cq_type::CQType;
use rocketmq_common::common::boundary_type::BoundaryType;
use rocketmq_common::common::broker::broker_config::BrokerConfig;
use rocketmq_common::common::message::message_ext_broker_inner::MessageExtBrokerInner;
use rocketmq_common::common::topic::TopicValidator;
use rocketmq_common::utils::queue_type_utils::QueueTypeUtils;
use rocketmq_rust::ArcMut;
use tracing::error;
use tracing::info;
use crate::base::dispatch_request::DispatchRequest;
use crate::base::message_store::MessageStore;
use crate::config::message_store_config::MessageStoreConfig;
use crate::message_store::local_file_message_store::LocalFileMessageStore;
use crate::queue::batch_consume_queue::BatchConsumeQueue;
use crate::queue::consume_queue::ConsumeQueueTrait;
use crate::queue::consume_queue_store::ConsumeQueueStoreTrait;
use crate::queue::queue_offset_operator::QueueOffsetOperator;
use crate::queue::single_consume_queue::ConsumeQueue;
use crate::queue::single_consume_queue::CQ_STORE_UNIT_SIZE;
use crate::queue::ArcConsumeQueue;
use crate::queue::ConsumeQueueTable;
use crate::queue::CqUnit;
use crate::store_path_config_helper::get_store_path_batch_consume_queue;
use crate::store_path_config_helper::get_store_path_consume_queue;
#[derive(Clone)]
pub struct ConsumeQueueStore {
inner: ArcMut<Inner>,
}
struct Inner {
pub(crate) message_store: Option<ArcMut<LocalFileMessageStore>>,
pub(crate) message_store_config: Arc<MessageStoreConfig>,
pub(crate) broker_config: Arc<BrokerConfig>,
pub(crate) queue_offset_operator: QueueOffsetOperator,
pub(crate) consume_queue_table: Arc<ConsumeQueueTable>,
}
impl Inner {
fn put_message_position_info_wrapper(&self, consume_queue: &mut dyn ConsumeQueueTrait, request: &DispatchRequest) {
consume_queue.put_message_position_info_wrapper(request)
}
}
impl ConsumeQueueStore {
pub fn clean_expired_sync(&self, min_commit_log_offset: i64) {
let mut queues_to_destroy = Vec::new();
let mut topics_to_remove = Vec::new();
{
let mut consume_queue_table = self.inner.consume_queue_table.lock();
let topics: Vec<CheetahString> = consume_queue_table.keys().cloned().collect();
for topic in topics {
if TopicValidator::is_system_topic(&topic) {
continue;
}
if let Some(queue_table) = consume_queue_table.get_mut(&topic) {
let queue_ids: Vec<i32> = queue_table.keys().cloned().collect();
let mut queues_to_remove = Vec::new();
for queue_id in queue_ids {
if let Some(consume_queue) = queue_table.get(&queue_id) {
let max_cl_offset_in_queue = consume_queue.get_max_physic_offset();
let message_total_in_queue = consume_queue.get_message_total_in_queue();
let queue_trimmed_to_empty =
message_total_in_queue <= 0 && consume_queue.get_min_offset_in_queue() > 0;
if max_cl_offset_in_queue == -1 {
tracing::warn!(
"maybe ConsumeQueue was created just now. topic={} queueId={} maxPhysicOffset={} \
minLogicOffset={}.",
consume_queue.get_topic(),
consume_queue.get_queue_id(),
consume_queue.get_max_physic_offset(),
consume_queue.get_min_logic_offset()
);
} else if max_cl_offset_in_queue < min_commit_log_offset || queue_trimmed_to_empty {
tracing::info!(
"cleanExpiredConsumerQueue: {} {} consumer queue destroyed, minCommitLogOffset: \
{} maxCLOffsetInConsumeQueue: {} messageTotalInQueue: {} minOffsetInQueue: {}",
topic,
queue_id,
min_commit_log_offset,
max_cl_offset_in_queue,
message_total_in_queue,
consume_queue.get_min_offset_in_queue()
);
queues_to_remove.push(queue_id);
}
}
}
for queue_id in queues_to_remove {
if let Some(consume_queue) = queue_table.remove(&queue_id) {
queues_to_destroy.push((topic.clone(), queue_id, consume_queue));
}
}
if queue_table.is_empty() {
topics_to_remove.push(topic.clone());
}
}
}
for topic in &topics_to_remove {
consume_queue_table.remove(topic);
tracing::info!("cleanExpiredConsumerQueue: {},topic destroyed", topic);
}
}
for (topic, queue_id, mut consume_queue) in queues_to_destroy {
self.inner.queue_offset_operator.remove(&topic, queue_id);
consume_queue.destroy();
}
}
fn find_consume_queue(&self, topic: &CheetahString, queue_id: i32) -> Option<ArcConsumeQueue> {
let table = self.inner.consume_queue_table.lock();
let queue_table = table.get(topic)?;
queue_table.get(&queue_id).cloned()
}
fn encode_cq_unit(cq_unit: &CqUnit) -> Bytes {
if !cq_unit.native_buffer.is_empty() {
return Bytes::copy_from_slice(&cq_unit.native_buffer);
}
let mut bytes = BytesMut::with_capacity(CQ_STORE_UNIT_SIZE as usize);
bytes.put_i64(cq_unit.pos);
bytes.put_i32(cq_unit.size);
bytes.put_i64(cq_unit.tags_code);
bytes.freeze()
}
#[inline]
pub fn new(message_store_config: Arc<MessageStoreConfig>, broker_config: Arc<BrokerConfig>) -> Self {
Self {
inner: ArcMut::new(Inner {
message_store: None,
message_store_config,
broker_config,
queue_offset_operator: Default::default(),
consume_queue_table: Arc::new(Default::default()),
}),
}
}
pub fn set_message_store(&mut self, message_store: ArcMut<LocalFileMessageStore>) {
self.inner.message_store = Some(message_store);
}
pub fn check_self(&self, consume_queue: &dyn ConsumeQueueTrait) {
let life_cycle = self.get_life_cycle(consume_queue.get_topic(), consume_queue.get_queue_id());
life_cycle.check_self();
}
}
#[allow(unused_variables)]
impl ConsumeQueueStoreTrait for ConsumeQueueStore {
fn start(&self) {
info!("consume queue store start");
}
fn load(&mut self) -> bool {
self.load_consume_queues(
get_store_path_consume_queue(self.inner.message_store_config.store_path_root_dir.as_str()).as_str(),
CQType::SimpleCQ,
) & self.load_consume_queues(
get_store_path_batch_consume_queue(self.inner.message_store_config.store_path_root_dir.as_str()).as_str(),
CQType::BatchCQ,
)
}
fn load_after_destroy(&self) -> bool {
true
}
async fn recover(&self) {
let mut mutex = self.inner.consume_queue_table.lock().clone();
for consume_queue_table in mutex.values_mut() {
for consume_queue in consume_queue_table.values() {
let queue_id = consume_queue.get_queue_id();
let topic = consume_queue.get_topic();
let mut file_queue_life_cycle = self.get_life_cycle(topic, queue_id);
file_queue_life_cycle.recover();
}
}
}
async fn recover_concurrently(&self) -> bool {
let mut count = 0;
for maps in self.inner.consume_queue_table.lock().values() {
count += maps.values().len();
}
let mut futures = Vec::with_capacity(count);
for maps in self.inner.consume_queue_table.lock().values() {
for logic in maps.values() {
let mut logic_clone = logic.clone();
let future = async move {
let mut ret = true;
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| logic_clone.recover())) {
Ok(_) => {}
Err(_) => {
ret = false;
error!(
"Exception occurs while recover consume queue concurrently, topic={}, queueId={}",
logic_clone.get_topic(),
logic_clone.get_queue_id()
);
}
}
ret
};
futures.push(future);
}
}
let results = join_all(futures).await;
for result in results {
if !result {
return false;
}
}
true
}
fn shutdown(&self) -> bool {
true
}
fn destroy(&self) {
let mutex = self.inner.consume_queue_table.lock().clone();
for consume_queue_table in mutex.values() {
for consume_queue in consume_queue_table.values() {
let queue_id = consume_queue.get_queue_id();
let topic = consume_queue.get_topic();
let mut file_queue_life_cycle = self.get_life_cycle(topic, queue_id);
file_queue_life_cycle.destroy();
}
}
}
fn destroy_queue(&self, consume_queue: &dyn ConsumeQueueTrait) {
let mut file_queue_life_cycle = self.get_life_cycle(consume_queue.get_topic(), consume_queue.get_queue_id());
file_queue_life_cycle.destroy();
}
fn flush(&self, consume_queue: &dyn ConsumeQueueTrait, flush_least_pages: i32) -> bool {
let file_queue_life_cycle = self.get_life_cycle(consume_queue.get_topic(), consume_queue.get_queue_id());
file_queue_life_cycle.flush(flush_least_pages)
}
async fn clean_expired(&self, min_commit_log_offset: i64) {
self.clean_expired_sync(min_commit_log_offset);
}
fn check_self(&self) {
let consume_queue_table = self.inner.consume_queue_table.lock().clone();
for consume_queue_table in consume_queue_table.values() {
for consume_queue in consume_queue_table.values() {
let consume_queue = &**consume_queue.as_ref();
ConsumeQueueStore::check_self(self, consume_queue);
}
}
}
fn delete_expired_file(&self, consume_queue: &dyn ConsumeQueueTrait, min_commit_log_pos: i64) -> i32 {
let file_queue_life_cycle = self.get_life_cycle(consume_queue.get_topic(), consume_queue.get_queue_id());
file_queue_life_cycle.delete_expired_file(min_commit_log_pos)
}
fn is_first_file_available(&self, consume_queue: &dyn ConsumeQueueTrait) -> bool {
let file_queue_life_cycle = self.get_life_cycle(consume_queue.get_topic(), consume_queue.get_queue_id());
file_queue_life_cycle.is_first_file_available()
}
fn is_first_file_exist(&self, consume_queue: &dyn ConsumeQueueTrait) -> bool {
let file_queue_life_cycle = self.get_life_cycle(consume_queue.get_topic(), consume_queue.get_queue_id());
file_queue_life_cycle.is_first_file_exist()
}
fn roll_next_file(&self, consume_queue: &dyn ConsumeQueueTrait, offset: i64) -> i64 {
let file_queue_life_cycle = self.get_life_cycle(consume_queue.get_topic(), consume_queue.get_queue_id());
file_queue_life_cycle.roll_next_file(offset)
}
fn truncate_dirty(&self, offset_to_truncate: i64) {
let cloned = self.inner.consume_queue_table.lock().clone();
for consume_queue_table in cloned.values() {
for logic in consume_queue_table.values() {
let topic = logic.get_topic();
let queue_id = logic.get_queue_id();
self.truncate_dirty_logic_files(topic, queue_id, offset_to_truncate);
}
}
}
fn put_message_position_info_wrapper_with_cq(
&self,
consume_queue: &mut dyn ConsumeQueueTrait,
request: &DispatchRequest,
) {
self.inner.put_message_position_info_wrapper(consume_queue, request);
}
fn put_message_position_info_wrapper(&self, request: &DispatchRequest) {
let mut cq = self.find_or_create_consume_queue(request.topic.as_ref(), request.queue_id);
self.put_message_position_info_wrapper_with_cq(cq.as_mut().as_mut(), request);
}
async fn range_query(&self, topic: &CheetahString, queue_id: i32, start_index: i64, num: i32) -> Vec<Bytes> {
let Some(consume_queue) = self.find_consume_queue(topic, queue_id) else {
return Vec::new();
};
let Some(iter) = consume_queue.iterate_from_with_count(start_index, num) else {
return Vec::new();
};
iter.take(num.max(0) as usize)
.map(|cq_unit| Self::encode_cq_unit(&cq_unit))
.collect()
}
async fn get(&self, topic: &CheetahString, queue_id: i32, start_index: i64) -> Bytes {
self.find_consume_queue(topic, queue_id)
.and_then(|consume_queue| consume_queue.get(start_index))
.map(|cq_unit| Self::encode_cq_unit(&cq_unit))
.unwrap_or_default()
}
fn get_consume_queue_table(&self) -> Arc<ConsumeQueueTable> {
self.inner.consume_queue_table.clone()
}
fn assign_queue_offset(&self, msg: &mut MessageExtBrokerInner) {
let consume_queue = self.find_or_create_consume_queue(msg.get_topic(), msg.queue_id());
consume_queue.assign_queue_offset(&self.inner.queue_offset_operator, msg);
}
fn increase_queue_offset(&self, msg: &MessageExtBrokerInner, message_num: i16) {
let consume_queue = self.find_or_create_consume_queue(msg.get_topic(), msg.queue_id());
consume_queue.increase_queue_offset(&self.inner.queue_offset_operator, msg, message_num);
}
fn increase_lmq_offset(&self, queue_key: &str, message_num: i16) {
let queue_key = CheetahString::from(queue_key);
self.inner
.queue_offset_operator
.increase_queue_offset(queue_key.clone(), message_num);
self.inner
.queue_offset_operator
.increase_lmq_offset(&queue_key, message_num);
}
fn get_lmq_queue_offset(&self, queue_key: &str) -> i64 {
let queue_key = CheetahString::from(queue_key);
let lmq_offset = self.inner.queue_offset_operator.get_lmq_offset(&queue_key);
if lmq_offset > 0 {
lmq_offset
} else {
self.inner.queue_offset_operator.current_queue_offset(&queue_key)
}
}
fn get_lmq_num(&self) -> i32 {
self.inner.queue_offset_operator.get_lmq_num()
}
fn is_lmq_exist(&self, lmq_topic: &str) -> bool {
self.inner.queue_offset_operator.is_lmq_exist(lmq_topic)
}
fn recover_offset_table(&mut self, min_phy_offset: i64) {
let mut cq_offset_table = HashMap::with_capacity(1024);
let mut bcq_offset_table = HashMap::with_capacity(1024);
for (topic, consume_queue_table) in self.inner.consume_queue_table.lock().iter_mut() {
for (queue_id, consume_queue) in consume_queue_table.iter() {
let key = CheetahString::from_string(format!(
"{}-{}",
consume_queue.get_topic(),
consume_queue.get_queue_id()
));
let max_offset_in_queue = consume_queue.get_max_offset_in_queue();
if consume_queue.get_cq_type() == CQType::SimpleCQ {
cq_offset_table.insert(key, max_offset_in_queue);
} else {
bcq_offset_table.insert(key, max_offset_in_queue);
}
self.correct_min_offset(&***consume_queue, min_phy_offset)
}
}
if self.inner.message_store_config.duplication_enable || self.inner.broker_config.enable_controller_mode {
self.inner
.queue_offset_operator
.set_lmq_topic_queue_table(cq_offset_table);
} else {
self.set_topic_queue_table(cq_offset_table);
self.set_batch_topic_queue_table(bcq_offset_table);
}
}
fn set_topic_queue_table(&mut self, topic_queue_table: HashMap<CheetahString, i64>) {
self.inner
.queue_offset_operator
.set_topic_queue_table(topic_queue_table.clone());
self.inner
.queue_offset_operator
.set_lmq_topic_queue_table(topic_queue_table);
}
fn remove_topic_queue_table(&mut self, topic: &CheetahString, queue_id: i32) {
self.inner.queue_offset_operator.remove(topic, queue_id);
}
fn get_topic_queue_table(&self) -> HashMap<CheetahString, i64> {
self.inner.queue_offset_operator.get_topic_queue_table()
}
fn get_max_phy_offset_in_consume_queue(&self, topic: &CheetahString, queue_id: i32) -> Option<i64> {
let mut max_physic_offset = -1i64;
for (topic, consume_queue_table) in self.inner.consume_queue_table.lock().iter() {
for (queue_id, consume_queue) in consume_queue_table.iter() {
let max_physic_offset_in_consume_queue = consume_queue.get_max_physic_offset();
if max_physic_offset_in_consume_queue > max_physic_offset {
max_physic_offset = max_physic_offset_in_consume_queue;
}
}
}
Some(max_physic_offset)
}
fn get_max_offset(&self, topic: &CheetahString, queue_id: i32) -> Option<i64> {
Some(
self.inner
.queue_offset_operator
.current_queue_offset(&format!("{topic}-{queue_id}").into()),
)
}
fn get_max_phy_offset_in_consume_queue_global(&self) -> i64 {
let mut max_physic_offset = -1i64;
for (topic, consume_queue_table) in self.inner.consume_queue_table.lock().iter() {
for (queue_id, consume_queue) in consume_queue_table.iter() {
let max_physic_offset_in_consume_queue = consume_queue.get_max_physic_offset();
if max_physic_offset_in_consume_queue > max_physic_offset {
max_physic_offset = max_physic_offset_in_consume_queue;
}
}
}
max_physic_offset
}
fn get_min_offset_in_queue(&self, topic: &CheetahString, queue_id: i32) -> i64 {
let queue = self.find_or_create_consume_queue(topic, queue_id);
queue.get_min_offset_in_queue()
}
fn get_max_offset_in_queue(&self, topic: &CheetahString, queue_id: i32) -> i64 {
let queue = self.find_or_create_consume_queue(topic, queue_id);
queue.get_max_offset_in_queue()
}
fn get_offset_in_queue_by_time(
&self,
topic: &CheetahString,
queue_id: i32,
timestamp: i64,
boundary_type: BoundaryType,
) -> i64 {
let logic = self.find_or_create_consume_queue(topic, queue_id);
let result_offset = logic.get_offset_in_queue_by_time_with_boundary(timestamp, boundary_type);
result_offset
.max(logic.get_min_offset_in_queue())
.min(logic.get_max_offset_in_queue())
}
fn find_or_create_consume_queue(&self, topic: &CheetahString, queue_id: i32) -> ArcConsumeQueue {
{
let consume_queue_table = self.inner.consume_queue_table.lock();
if let Some(topic_map) = consume_queue_table.get(topic) {
if let Some(queue) = topic_map.get(&queue_id) {
return queue.clone();
}
}
}
let message_store = self
.inner
.message_store
.as_ref()
.expect("MessageStore must be set before creating consume queues");
let topic_config = message_store.get_topic_config(topic);
let cq_type = QueueTypeUtils::get_cq_type_arc_mut(topic_config.as_ref());
let new_queue: ArcConsumeQueue = match cq_type {
CQType::SimpleCQ | CQType::RocksDBCQ => {
let consume_queue = ConsumeQueue::new(
topic.clone(),
queue_id,
CheetahString::from_string(get_store_path_consume_queue(
self.inner.message_store_config.store_path_root_dir.as_str(),
)),
self.inner.message_store_config.get_mapped_file_size_consume_queue(),
self.inner.message_store.clone().unwrap(),
);
ArcMut::new(Box::new(consume_queue))
}
CQType::BatchCQ => {
let consume_queue = BatchConsumeQueue::new(
topic.clone(),
queue_id,
CheetahString::from_string(get_store_path_batch_consume_queue(
self.inner.message_store_config.store_path_root_dir.as_str(),
)),
self.inner.message_store_config.mapper_file_size_batch_consume_queue,
None,
self.inner.message_store_config.clone(),
);
ArcMut::new(Box::new(consume_queue))
}
};
let mut consume_queue_table = self.inner.consume_queue_table.lock();
let topic_map = consume_queue_table.entry(topic.clone()).or_default();
topic_map.entry(queue_id).or_insert(new_queue).clone()
}
fn find_consume_queue_map(&self, topic: &CheetahString) -> Option<HashMap<i32, ArcConsumeQueue>> {
self.inner.consume_queue_table.lock().get(topic).cloned()
}
fn get_total_size(&self) -> i64 {
let mut total_size = 0;
for consume_queue_table in self.inner.consume_queue_table.lock().values() {
for consume_queue in consume_queue_table.values() {
total_size += consume_queue.get_total_size();
}
}
total_size
}
fn get_store_time(&self, cq_unit: &CqUnit) -> i64 {
match self.inner.message_store.as_ref() {
Some(ms) => ms.get_commit_log().pickup_store_timestamp(cq_unit.pos, cq_unit.size),
None => {
error!("Message store is not set in ConsumeQueueStore");
-1
}
}
}
fn as_any(&self) -> &dyn Any {
self
}
fn as_any_mut(&mut self) -> &mut dyn Any {
self
}
}
impl ConsumeQueueStore {
#[inline]
pub fn correct_min_offset(&self, consume_queue: &dyn ConsumeQueueTrait, min_commit_log_offset: i64) {
consume_queue.correct_min_offset(min_commit_log_offset)
}
#[inline]
pub fn set_batch_topic_queue_table(&self, batch_topic_queue_table: HashMap<CheetahString, i64>) {
self.inner
.queue_offset_operator
.set_batch_topic_queue_table(batch_topic_queue_table)
}
fn load_consume_queues(&mut self, store_path: &str, cq_type: CQType) -> bool {
let dir_logic = Path::new(store_path);
if !dir_logic.exists() || !dir_logic.is_dir() {
info!("Directory {} doesn't exist or is not a directory", store_path);
return true; }
match fs::read_dir(dir_logic) {
Ok(topic_entries) => {
for topic_dir in topic_entries.flatten() {
let topic_path = topic_dir.path();
if !topic_path.is_dir() {
continue;
}
let topic = match topic_path.file_name().and_then(|n| n.to_str()) {
Some(name) => CheetahString::from_string(name.to_string()),
None => continue,
};
match fs::read_dir(topic_path) {
Ok(queue_id_entries) => {
for queue_id_dir in queue_id_entries.flatten() {
if !queue_id_dir.path().is_dir() {
continue;
}
let os_string = queue_id_dir.file_name();
let queue_id_op = os_string.to_str();
let queue_id_str = match queue_id_op {
Some(name) => name,
None => continue,
};
let queue_id = match i32::from_str(queue_id_str) {
Ok(id) => id,
Err(_) => {
continue;
}
};
if !self.queue_type_should_be(&topic, cq_type) {
return false;
}
let logic =
self.create_consume_queue_by_type(&topic, queue_id, cq_type, store_path.into());
self.put_consume_queue(topic.clone(), queue_id, logic);
if !self.load_logic(&topic, queue_id) {
return false;
}
}
}
Err(e) => {
error!("Failed to read queue ID directories for topic {}: {}", topic, e);
return false;
}
}
}
}
Err(e) => {
error!("Failed to read topic directories: {e}");
return false;
}
}
info!("load {cq_type} all over, OK");
true
}
#[inline]
fn load_logic(&mut self, topic: &CheetahString, queue_id: i32) -> bool {
let mut file_queue_life_cycle = self.get_life_cycle(topic, queue_id);
file_queue_life_cycle.load()
}
#[inline]
fn put_consume_queue(
&self,
topic: CheetahString,
queue_id: i32,
consume_queue: ArcMut<Box<dyn ConsumeQueueTrait>>,
) {
let mut consume_queue_table = self.inner.consume_queue_table.lock();
let topic_table = consume_queue_table.entry(topic).or_default();
topic_table.insert(queue_id, consume_queue);
}
#[inline]
fn queue_type_should_be(&self, topic: &CheetahString, cq_type: CQType) -> bool {
let Some(message_store) = self.inner.message_store.as_ref() else {
error!("MessageStore must be set before loading consume queues for topic {topic}");
return false;
};
let topic_config = message_store.get_topic_config(topic);
let act = QueueTypeUtils::get_cq_type_arc_mut(topic_config.as_ref());
if self.is_rocksdb_cq_compat(act, cq_type) {
return true;
}
if act != cq_type {
error!("The queue type of topic: {topic} should be {cq_type:?}, but is {act:?}");
return false;
}
true
}
#[inline]
fn is_rocksdb_cq_compat(&self, actual: CQType, expected: CQType) -> bool {
self.inner.message_store_config.is_enable_rocksdb_store()
&& actual == CQType::RocksDBCQ
&& expected == CQType::SimpleCQ
}
#[inline]
fn truncate_dirty_logic_files(&self, topic: &CheetahString, queue_id: i32, phy_offset: i64) {
let mut file_queue_life_cycle = self.get_life_cycle(topic, queue_id);
file_queue_life_cycle.truncate_dirty_logic_files(phy_offset);
}
#[inline]
fn create_consume_queue_by_type(
&self,
topic: &CheetahString,
queue_id: i32,
cq_type: CQType,
store_path: CheetahString,
) -> ArcMut<Box<dyn ConsumeQueueTrait>> {
match cq_type {
CQType::SimpleCQ | CQType::RocksDBCQ => {
let consume_queue = ConsumeQueue::new(
topic.clone(),
queue_id,
store_path,
self.inner.message_store_config.get_mapped_file_size_consume_queue(),
self.inner.message_store.clone().unwrap(),
);
ArcMut::new(Box::new(consume_queue))
}
CQType::BatchCQ => {
let consume_queue = BatchConsumeQueue::new(
topic.clone(),
queue_id,
store_path,
self.inner.message_store_config.mapper_file_size_batch_consume_queue,
None,
self.inner.message_store_config.clone(),
);
ArcMut::new(Box::new(consume_queue))
}
}
}
#[inline]
fn get_life_cycle(&self, topic: &CheetahString, queue_id: i32) -> ArcConsumeQueue {
self.find_or_create_consume_queue(topic, queue_id)
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use bytes::Buf;
use cheetah_string::CheetahString;
use dashmap::DashMap;
use rocketmq_common::common::attribute::Attribute;
use rocketmq_common::common::broker::broker_config::BrokerConfig;
use rocketmq_common::common::config::TopicConfig;
use rocketmq_common::TopicAttributes::TopicAttributes;
use rocketmq_rust::ArcMut;
use tempfile::tempdir;
use super::*;
use crate::base::store_enum::StoreType;
use crate::queue::batch_consume_queue;
#[test]
fn load_consume_queues_returns_false_when_topic_queue_type_mismatches_directory() {
let root = tempdir().expect("tempdir");
let topic = CheetahString::from_static_str("MismatchTopic");
let message_store_config = Arc::new(MessageStoreConfig {
store_path_root_dir: root.path().to_string_lossy().to_string().into(),
..MessageStoreConfig::default()
});
let broker_config = Arc::new(BrokerConfig::default());
let topic_config_table = Arc::new(DashMap::<CheetahString, ArcMut<TopicConfig>>::new());
topic_config_table.insert(topic.clone(), ArcMut::new(TopicConfig::new(topic.clone())));
let mut message_store = ArcMut::new(LocalFileMessageStore::new(
message_store_config.clone(),
broker_config.clone(),
topic_config_table,
None,
false,
));
let message_store_clone = message_store.clone();
message_store.set_message_store_arc(message_store_clone);
let batch_store_path = get_store_path_batch_consume_queue(message_store_config.store_path_root_dir.as_str());
fs::create_dir_all(Path::new(&batch_store_path).join(topic.as_str()).join("0"))
.expect("batch consume queue directory");
let mut store = ConsumeQueueStore::new(message_store_config, broker_config);
store.set_message_store(message_store);
assert!(!store.load_consume_queues(&batch_store_path, CQType::BatchCQ));
}
#[test]
fn rocksdb_cq_loads_from_simple_cq_path_in_compat_mode() {
let root = tempdir().expect("tempdir");
let topic = CheetahString::from_static_str("RocksCompatTopic");
let message_store_config = Arc::new(MessageStoreConfig {
store_path_root_dir: root.path().to_string_lossy().to_string().into(),
store_type: StoreType::RocksDB,
..MessageStoreConfig::default()
});
let broker_config = Arc::new(BrokerConfig::default());
let topic_config_table = Arc::new(DashMap::<CheetahString, ArcMut<TopicConfig>>::new());
let mut topic_config = TopicConfig::new(topic.clone());
topic_config.attributes.insert(
TopicAttributes::queue_type_attribute().name().clone(),
CQType::RocksDBCQ.to_string().into(),
);
topic_config_table.insert(topic.clone(), ArcMut::new(topic_config));
let mut message_store = ArcMut::new(LocalFileMessageStore::new(
message_store_config.clone(),
broker_config.clone(),
topic_config_table,
None,
false,
));
let message_store_clone = message_store.clone();
message_store.set_message_store_arc(message_store_clone);
let simple_store_path = get_store_path_consume_queue(message_store_config.store_path_root_dir.as_str());
fs::create_dir_all(Path::new(&simple_store_path).join(topic.as_str()).join("0"))
.expect("simple consume queue directory");
let mut store = ConsumeQueueStore::new(message_store_config, broker_config);
store.set_message_store(message_store);
assert!(store.load_consume_queues(&simple_store_path, CQType::SimpleCQ));
let queue = store
.find_consume_queue(&topic, 0)
.expect("rocksdb compatibility queue should load");
assert_eq!(queue.get_cq_type(), CQType::SimpleCQ);
}
#[test]
fn create_consume_queue_by_type_maps_rocksdb_cq_to_simple_cq() {
let root = tempdir().expect("tempdir");
let topic = CheetahString::from_static_str("RocksCreateTopic");
let message_store_config = Arc::new(MessageStoreConfig {
store_path_root_dir: root.path().to_string_lossy().to_string().into(),
store_type: StoreType::RocksDB,
..MessageStoreConfig::default()
});
let broker_config = Arc::new(BrokerConfig::default());
let topic_config_table = Arc::new(DashMap::<CheetahString, ArcMut<TopicConfig>>::new());
let mut message_store = ArcMut::new(LocalFileMessageStore::new(
message_store_config.clone(),
broker_config.clone(),
topic_config_table,
None,
false,
));
let message_store_clone = message_store.clone();
message_store.set_message_store_arc(message_store_clone);
let mut store = ConsumeQueueStore::new(message_store_config, broker_config);
store.set_message_store(message_store);
let queue = store.create_consume_queue_by_type(
&topic,
0,
CQType::RocksDBCQ,
CheetahString::from_string(root.path().join("compat-cq").to_string_lossy().to_string()),
);
assert_eq!(queue.get_cq_type(), CQType::SimpleCQ);
}
#[test]
fn batch_consume_queue_store_get_returns_full_batch_unit() {
let message_store_config = Arc::new(MessageStoreConfig::default());
let broker_config = Arc::new(BrokerConfig::default());
let store = ConsumeQueueStore::new(message_store_config.clone(), broker_config);
let root = tempdir().expect("tempdir");
let store_path = CheetahString::from_string(root.path().join("batch-cq").to_string_lossy().to_string());
let topic = CheetahString::from_static_str("BatchTopic");
let queue_id = 1;
let mut batch_queue = BatchConsumeQueue::new(
topic.clone(),
queue_id,
store_path,
(batch_consume_queue::CQ_STORE_UNIT_SIZE * 2) as usize,
None,
message_store_config,
);
batch_queue.put_message_position_info_wrapper(&DispatchRequest {
topic: topic.clone(),
queue_id,
commit_log_offset: 100,
msg_size: 32,
tags_code: 7,
store_timestamp: 1_000,
msg_base_offset: 0,
batch_size: 3,
success: true,
..DispatchRequest::default()
});
store
.inner
.consume_queue_table
.lock()
.entry(topic.clone())
.or_default()
.insert(queue_id, ArcMut::new(Box::new(batch_queue)));
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
let bytes = runtime.block_on(store.get(&topic, queue_id, 0));
assert_eq!(bytes.len(), batch_consume_queue::CQ_STORE_UNIT_SIZE as usize);
let mut bytes = bytes;
assert_eq!(bytes.get_i64(), 100);
assert_eq!(bytes.get_i32(), 32);
assert_eq!(bytes.get_i64(), 7);
assert_eq!(bytes.get_i64(), 1_000);
assert_eq!(bytes.get_i64(), 0);
assert_eq!(bytes.get_i16(), 3);
}
}