1use std::sync::Arc;
16
17use api::helper::to_pb_time_unit;
18use api::v1::region::region_request::Body as RegionRequestBody;
19use api::v1::region::{
20 BuildIndexRequest, CompactRequest, CompactionTimeRange, FlushRequest, RegionRequestHeader,
21 TruncateRequest, Unflushed, build_index_request, truncate_request,
22};
23use catalog::CatalogManagerRef;
24use common_catalog::build_db_string;
25use common_catalog::consts::METRIC_ENGINE;
26use common_meta::node_manager::{AffectedRows, NodeManagerRef};
27use common_meta::peer::Peer;
28use common_telemetry::tracing_context::TracingContext;
29use common_telemetry::{debug, error, info};
30use common_time::range::TimestampRange;
31use common_time::timestamp::TimeUnit as TimestampUnit;
32use futures_util::future;
33use partition::cache::PhysicalPartitionInfo;
34use partition::manager::PartitionRuleManagerRef;
35use session::context::QueryContextRef;
36use snafu::prelude::*;
37use store_api::storage::RegionId;
38use table::requests::{BuildIndexTableRequest, CompactTableRequest, FlushTableRequest};
39use table::table_name::TableName;
40
41use crate::error::{
42 CatalogSnafu, FindRegionLeaderSnafu, FindTablePartitionRuleSnafu, JoinTaskSnafu,
43 NotSupportedSnafu, RequestRegionSnafu, Result, TableNotFoundSnafu,
44 UnsupportedRegionRequestSnafu,
45};
46use crate::region_req_factory::RegionRequestFactory;
47
48pub struct Requester {
50 catalog_manager: CatalogManagerRef,
51 partition_manager: PartitionRuleManagerRef,
52 node_manager: NodeManagerRef,
53}
54
55pub type RequesterRef = Arc<Requester>;
56
57impl Requester {
58 pub fn new(
59 catalog_manager: CatalogManagerRef,
60 partition_manager: PartitionRuleManagerRef,
61 node_manager: NodeManagerRef,
62 ) -> Self {
63 Self {
64 catalog_manager,
65 partition_manager,
66 node_manager,
67 }
68 }
69
70 pub async fn handle_table_flush(
72 &self,
73 request: FlushTableRequest,
74 ctx: QueryContextRef,
75 ) -> Result<AffectedRows> {
76 let partitions = &self
77 .get_table_partition_info(
78 &request.catalog_name,
79 &request.schema_name,
80 &request.table_name,
81 )
82 .await?
83 .partitions;
84
85 let requests = partitions
86 .iter()
87 .map(|partition| {
88 RegionRequestBody::Flush(FlushRequest {
89 region_id: partition.id.into(),
90 })
91 })
92 .collect();
93
94 info!("Handle table manual flush request: {:?}", request);
95
96 self.do_request(
97 requests,
98 Some(build_db_string(&request.catalog_name, &request.schema_name)),
99 &ctx,
100 )
101 .await
102 }
103
104 pub async fn handle_table_build_index(
106 &self,
107 request: BuildIndexTableRequest,
108 ctx: QueryContextRef,
109 ) -> Result<AffectedRows> {
110 if matches!(
111 request.options,
112 Some(build_index_request::Options::SeriesIndex(_))
113 ) {
114 let table = self
115 .catalog_manager
116 .table(
117 &request.catalog_name,
118 &request.schema_name,
119 &request.table_name,
120 None,
121 )
122 .await
123 .context(CatalogSnafu)?;
124 let table = table.with_context(|| TableNotFoundSnafu {
125 table_name: common_catalog::format_full_table_name(
126 &request.catalog_name,
127 &request.schema_name,
128 &request.table_name,
129 ),
130 })?;
131 let info = table.table_info();
132 ensure_build_series_index_supported(&info.meta.engine, info.is_physical_table())?;
133 }
134 let partitions = &self
135 .get_table_partition_info(
136 &request.catalog_name,
137 &request.schema_name,
138 &request.table_name,
139 )
140 .await?
141 .partitions;
142
143 let requests = partitions
144 .iter()
145 .map(|partition| {
146 RegionRequestBody::BuildIndex(BuildIndexRequest {
147 region_id: partition.id.into(),
148 options: request.options,
149 })
150 })
151 .collect();
152
153 info!(
154 "Handle table manual build index for table {}",
155 request.table_name
156 );
157 debug!("Request details: {:?}", request);
158
159 self.do_request(
160 requests,
161 Some(build_db_string(&request.catalog_name, &request.schema_name)),
162 &ctx,
163 )
164 .await
165 }
166
167 pub async fn handle_table_compaction(
169 &self,
170 request: CompactTableRequest,
171 ctx: QueryContextRef,
172 ) -> Result<AffectedRows> {
173 let partitions = &self
174 .get_table_partition_info(
175 &request.catalog_name,
176 &request.schema_name,
177 &request.table_name,
178 )
179 .await?
180 .partitions;
181
182 let time_range = request
183 .time_range
184 .map(to_pb_compaction_time_range)
185 .transpose()?;
186 let requests = partitions
187 .iter()
188 .map(|partition| {
189 RegionRequestBody::Compact(CompactRequest {
190 region_id: partition.id.into(),
191 parallelism: request.parallelism,
192 options: Some(request.compact_options),
193 time_range,
194 })
195 })
196 .collect();
197
198 info!("Handle table manual compaction request: {:?}", request);
199
200 self.do_request(
201 requests,
202 Some(build_db_string(&request.catalog_name, &request.schema_name)),
203 &ctx,
204 )
205 .await
206 }
207
208 pub async fn handle_region_flush(
210 &self,
211 region_id: RegionId,
212 ctx: QueryContextRef,
213 ) -> Result<AffectedRows> {
214 let request = RegionRequestBody::Flush(FlushRequest {
215 region_id: region_id.into(),
216 });
217
218 info!("Handle region manual flush request: {region_id}");
219 self.do_request(vec![request], None, &ctx).await
220 }
221
222 pub async fn handle_region_compaction(
224 &self,
225 region_id: RegionId,
226 ctx: QueryContextRef,
227 ) -> Result<AffectedRows> {
228 let request = RegionRequestBody::Compact(CompactRequest {
229 region_id: region_id.into(),
230 parallelism: 1,
231 options: None, time_range: None,
233 });
234
235 info!("Handle region manual compaction request: {region_id}");
236 self.do_request(vec![request], None, &ctx).await
237 }
238
239 pub async fn handle_discard_unflushed_data(
241 &self,
242 region_id: RegionId,
243 ctx: QueryContextRef,
244 ) -> Result<AffectedRows> {
245 let request = RegionRequestBody::Truncate(TruncateRequest {
246 region_id: region_id.into(),
247 kind: Some(truncate_request::Kind::Unflushed(Unflushed {})),
248 });
249
250 info!("Handle region discard unflushed data request: {region_id}");
251 self.do_request(vec![request], None, &ctx).await
252 }
253
254 pub async fn handle_discard_unflushed_data_by_table(
256 &self,
257 table_name: TableName,
258 ctx: QueryContextRef,
259 ) -> Result<AffectedRows> {
260 let table = self
261 .catalog_manager
262 .table(
263 &table_name.catalog_name,
264 &table_name.schema_name,
265 &table_name.table_name,
266 None,
267 )
268 .await
269 .context(CatalogSnafu)?;
270 let table = table.with_context(|| TableNotFoundSnafu {
271 table_name: table_name.to_string(),
272 })?;
273 let table_info = table.table_info();
274 ensure_discard_unflushed_supported(
275 &table_info.meta.engine,
276 table_info.is_physical_table(),
277 )?;
278
279 let partitions = &self
280 .partition_manager
281 .find_physical_partition_info(table_info.ident.table_id)
282 .await
283 .with_context(|_| FindTablePartitionRuleSnafu {
284 table_name: table_name.to_string(),
285 })?
286 .partitions;
287 let requests = partitions
288 .iter()
289 .map(|partition| {
290 RegionRequestBody::Truncate(TruncateRequest {
291 region_id: partition.id.into(),
292 kind: Some(truncate_request::Kind::Unflushed(Unflushed {})),
293 })
294 })
295 .collect();
296
297 info!("Handle table discard unflushed data request: {table_name}");
298 self.do_request(
299 requests,
300 Some(build_db_string(
301 &table_name.catalog_name,
302 &table_name.schema_name,
303 )),
304 &ctx,
305 )
306 .await
307 }
308}
309
310fn to_pb_compaction_time_range(range: TimestampRange) -> Result<CompactionTimeRange> {
311 let (Some(start), Some(end)) = (*range.start(), *range.end()) else {
312 return crate::error::InvalidTimestampRangeSnafu {
313 start: format!("{:?}", range.start()),
314 end: format!("{:?}", range.end()),
315 }
316 .fail();
317 };
318 let original_start = start;
319 let original_end = end;
320 let start = start.convert_to(TimestampUnit::Second).with_context(|| {
321 crate::error::InvalidTimestampRangeSnafu {
322 start: format!("{original_start:?}"),
323 end: format!("{original_end:?}"),
324 }
325 })?;
326 let end = end
327 .convert_to_ceil(TimestampUnit::Second)
328 .with_context(|| crate::error::InvalidTimestampRangeSnafu {
329 start: format!("{original_start:?}"),
330 end: format!("{original_end:?}"),
331 })?;
332 ensure!(
333 start < end,
334 crate::error::InvalidTimestampRangeSnafu {
335 start: format!("{original_start:?}"),
336 end: format!("{original_end:?}"),
337 }
338 );
339
340 Ok(CompactionTimeRange {
341 start: start.value(),
342 end: end.value(),
343 time_unit: to_pb_time_unit(TimestampUnit::Second) as i32,
344 })
345}
346
347impl Requester {
348 async fn do_request(
349 &self,
350 requests: Vec<RegionRequestBody>,
351 db_string: Option<String>,
352 ctx: &QueryContextRef,
353 ) -> Result<AffectedRows> {
354 let request_factory = RegionRequestFactory::new(RegionRequestHeader {
355 tracing_context: TracingContext::from_current_span().to_w3c(),
356 dbname: db_string.unwrap_or_else(|| ctx.get_db_string()),
357 ..Default::default()
358 });
359
360 let tasks = requests.into_iter().map(|req_body| {
361 let request = request_factory.build_request(req_body.clone());
362 let partition_manager = self.partition_manager.clone();
363 let node_manager = self.node_manager.clone();
364 common_runtime::spawn_global(async move {
365 let peer =
366 Self::find_region_leader_by_request(partition_manager, &req_body).await?;
367 node_manager
368 .datanode(&peer)
369 .await
370 .handle(request)
371 .await
372 .context(RequestRegionSnafu)
373 })
374 });
375 let results = future::try_join_all(tasks).await.context(JoinTaskSnafu)?;
376
377 let affected_rows = results
378 .into_iter()
379 .map(|resp| resp.map(|r| r.affected_rows))
380 .sum::<Result<AffectedRows>>()?;
381
382 Ok(affected_rows)
383 }
384
385 async fn find_region_leader_by_request(
386 partition_manager: PartitionRuleManagerRef,
387 req: &RegionRequestBody,
388 ) -> Result<Peer> {
389 let region_id = match req {
390 RegionRequestBody::Flush(req) => req.region_id,
391 RegionRequestBody::Compact(req) => req.region_id,
392 RegionRequestBody::BuildIndex(req) => req.region_id,
393 RegionRequestBody::Truncate(req) => req.region_id,
394 _ => {
395 error!("Unsupported region request: {:?}", req);
396 return UnsupportedRegionRequestSnafu {}.fail();
397 }
398 };
399
400 partition_manager
401 .find_region_leader(region_id.into())
402 .await
403 .context(FindRegionLeaderSnafu)
404 }
405
406 async fn get_table_partition_info(
407 &self,
408 catalog: &str,
409 schema: &str,
410 table_name: &str,
411 ) -> Result<Arc<PhysicalPartitionInfo>> {
412 let table = self
413 .catalog_manager
414 .table(catalog, schema, table_name, None)
415 .await
416 .context(CatalogSnafu)?;
417
418 let table = table.with_context(|| TableNotFoundSnafu {
419 table_name: common_catalog::format_full_table_name(catalog, schema, table_name),
420 })?;
421 let table_info = table.table_info();
422
423 self.partition_manager
424 .find_physical_partition_info(table_info.ident.table_id)
425 .await
426 .with_context(|_| FindTablePartitionRuleSnafu {
427 table_name: common_catalog::format_full_table_name(catalog, schema, table_name),
428 })
429 }
430}
431
432fn ensure_build_series_index_supported(engine: &str, is_physical_table: bool) -> Result<()> {
433 ensure!(
434 engine == METRIC_ENGINE && is_physical_table,
435 NotSupportedSnafu {
436 feat: "building series indexes requires a physical metric table",
437 }
438 );
439 Ok(())
440}
441
442fn ensure_discard_unflushed_supported(engine: &str, is_physical_table: bool) -> Result<()> {
443 ensure!(
444 engine != METRIC_ENGINE || is_physical_table,
445 NotSupportedSnafu {
446 feat: "discarding unflushed data from a Metric Engine logical table"
447 }
448 );
449 Ok(())
450}
451
452#[cfg(test)]
453mod tests {
454 use api::v1::TimeUnit;
455 use common_time::Timestamp;
456 use common_time::range::TimestampRange;
457
458 use super::*;
459
460 #[test]
461 fn test_build_series_index_requires_physical_metric_table() {
462 ensure_build_series_index_supported(METRIC_ENGINE, true).unwrap();
463 assert!(ensure_build_series_index_supported(METRIC_ENGINE, false).is_err());
464 assert!(ensure_build_series_index_supported("mito", false).is_err());
465 assert!(ensure_build_series_index_supported("mito", true).is_err());
466 }
467
468 #[test]
469 fn test_to_pb_compaction_time_range_normalizes_mixed_units() {
470 let range = TimestampRange::new(
471 Timestamp::new_millisecond(1_500),
472 Timestamp::new_microsecond(2_500_000),
473 )
474 .unwrap();
475
476 let pb_range = to_pb_compaction_time_range(range).unwrap();
477 assert_eq!(1, pb_range.start);
478 assert_eq!(3, pb_range.end);
479 assert_eq!(TimeUnit::Second as i32, pb_range.time_unit);
480 }
481
482 #[test]
483 fn test_discard_unflushed_rejects_metric_logical_table() {
484 let error = ensure_discard_unflushed_supported(METRIC_ENGINE, false).unwrap_err();
485 assert!(matches!(error, crate::error::Error::NotSupported { .. }));
486
487 ensure_discard_unflushed_supported(METRIC_ENGINE, true).unwrap();
488 ensure_discard_unflushed_supported("mito", false).unwrap();
489 }
490
491 #[test]
492 fn test_to_pb_compaction_time_range() {
493 let range = TimestampRange::new(
494 Timestamp::new_microsecond(1_000),
495 Timestamp::new_microsecond(2_000),
496 )
497 .unwrap();
498
499 let pb_range = to_pb_compaction_time_range(range).unwrap();
500 assert_eq!(0, pb_range.start);
501 assert_eq!(1, pb_range.end);
502 assert_eq!(TimeUnit::Second as i32, pb_range.time_unit);
503 }
504}