-
Notifications
You must be signed in to change notification settings - Fork 273
Expand file tree
/
Copy pathmetadata.rs
More file actions
114 lines (109 loc) Β· 4.11 KB
/
Copy pathmetadata.rs
File metadata and controls
114 lines (109 loc) Β· 4.11 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
// Copyright 2023 The RocketMQ Rust Authors
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use super::*;
pub(super) struct BrokerMetadata {
pub(super) configuration_error: Option<String>,
#[cfg(feature = "rocksdb_store")]
pub(super) rocksdb_config_managers: Option<BrokerRocksDbConfigManagers>,
}
impl BrokerMetadata {
pub(super) fn new(
configuration_error: Option<String>,
#[cfg(feature = "rocksdb_store")] rocksdb_config_managers: Option<BrokerRocksDbConfigManagers>,
) -> Self {
Self {
configuration_error,
#[cfg(feature = "rocksdb_store")]
rocksdb_config_managers,
}
}
}
impl BrokerRuntime {
/// Load the original configuration data from the corresponding configuration files
/// located under the `${HOME}\config` directory.
///
/// This function initializes broker metadata by loading several manager components
/// in sequence:
/// - Topic configuration manager
/// - Topic queue mapping manager
/// - Consumer offset manager
/// - Subscription group manager
/// - Consumer filter manager
/// - Consumer order information manager
///
/// The loaders are invoked in order and combined using logical AND. If all loaders
/// return `true`, the function returns `true`. If any loader fails (returns `false`),
/// the whole initialization is considered failed and the function returns `false`.
pub(super) async fn initialize_metadata(&self) -> Result<(), BrokerStartupError> {
info!("======Starting initialize metadata========");
if let Some(Err(error)) = self.composition.state.metadata_io.as_ref() {
return Err(BrokerStartupError::initialization_source(
"metadata_io_actor",
error.clone(),
));
}
match self.composition.state.topic_config_coordinator().load().await {
Ok(true) => {}
Ok(false) => {
return Err(BrokerStartupError::metadata_load("topic_config"));
}
Err(error) => {
return Err(BrokerStartupError::initialization_source("topic_config", error));
}
}
for (component, loaded) in [
(
"topic_queue_mapping",
self.composition.state.topic_queue_mapping_manager().load(),
),
(
"consumer_offset",
self.composition.state.consumer_offset_manager().load(),
),
(
"subscription_group",
self.composition.state.subscription_group_manager().load(),
),
(
"consumer_filter",
self.composition.state.consumer_filter_manager().load(),
),
(
"consumer_order_info",
self.composition.state.consumer_order_info_manager().load(),
),
] {
if !loaded {
return Err(BrokerStartupError::metadata_load(component));
}
}
Ok(())
}
pub(super) async fn update_namesrv_addr(&mut self) {
self.composition.state.update_namesrv_addr_inner().await;
}
/// Register broker to name remoting_server
pub(crate) async fn register_broker_all(
&mut self,
check_order_config: bool,
oneway: bool,
force_register: bool,
) -> Result<BrokerRegistrationStatus, BrokerRegistrationError> {
self.composition
.state
.build_registration_runtime()
.register_broker_all(check_order_config, oneway, force_register)
.await
}
}