use rocketmq_common::common::sys_flag::message_sys_flag::MessageSysFlag;
use crate::base::commit_log_dispatcher::CommitLogDispatcher;
use crate::base::dispatch_request::DispatchRequest;
use crate::queue::local_file_consume_queue_store::ConsumeQueueStore;
use crate::queue::ConsumeQueueStoreTrait;
#[derive(Clone)]
pub struct CommitLogDispatcherBuildConsumeQueue {
consume_queue_store: ConsumeQueueStore,
}
impl CommitLogDispatcherBuildConsumeQueue {
pub fn new(consume_queue_store: ConsumeQueueStore) -> Self {
Self {
consume_queue_store,
}
}
}
impl CommitLogDispatcher for CommitLogDispatcherBuildConsumeQueue {
fn dispatch(&self, dispatch_request: &DispatchRequest) {
let tran_type = MessageSysFlag::get_transaction_value(dispatch_request.sys_flag);
match tran_type {
MessageSysFlag::TRANSACTION_NOT_TYPE | MessageSysFlag::TRANSACTION_COMMIT_TYPE => {
self.consume_queue_store
.put_message_position_info_wrapper(dispatch_request);
}
_ => {}
}
}
}