Skip to main content

common_function/admin/
build_index_table.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use api::v1::region::build_index_request;
16use arrow::datatypes::DataType as ArrowDataType;
17use common_error::ext::BoxedError;
18use common_macro::admin_fn;
19use common_query::error::{
20    InvalidFuncArgsSnafu, MissingTableMutationHandlerSnafu, Result, TableMutationSnafu,
21    UnsupportedInputDataTypeSnafu,
22};
23use datafusion_expr::{Signature, Volatility};
24use datatypes::prelude::*;
25use session::context::QueryContextRef;
26use session::table_name::table_name_to_full_name;
27use snafu::{ResultExt, ensure};
28use table::requests::BuildIndexTableRequest;
29
30use crate::handlers::TableMutationHandlerRef;
31
32#[admin_fn(
33    name = BuildIndexFunction,
34    display_name = build_index,
35    sig_fn = build_index_signature,
36    ret = uint64
37)]
38pub(crate) async fn build_index(
39    table_mutation_handler: &TableMutationHandlerRef,
40    query_ctx: &QueryContextRef,
41    params: &[ValueRef<'_>],
42) -> Result<Value> {
43    build_index_impl(
44        table_mutation_handler,
45        query_ctx,
46        params,
47        "build_index",
48        None,
49    )
50    .await
51}
52
53/// Reconciles series indexes on every physical data-region leader and waits for publication.
54#[admin_fn(
55    name = BuildSeriesIndexFunction,
56    display_name = build_series_index,
57    sig_fn = build_index_signature,
58    ret = uint64
59)]
60pub(crate) async fn build_series_index(
61    table_mutation_handler: &TableMutationHandlerRef,
62    query_ctx: &QueryContextRef,
63    params: &[ValueRef<'_>],
64) -> Result<Value> {
65    build_index_impl(
66        table_mutation_handler,
67        query_ctx,
68        params,
69        "build_series_index",
70        Some(build_index_request::Options::SeriesIndex(Default::default())),
71    )
72    .await
73}
74
75async fn build_index_impl(
76    table_mutation_handler: &TableMutationHandlerRef,
77    query_ctx: &QueryContextRef,
78    params: &[ValueRef<'_>],
79    function: &str,
80    options: Option<build_index_request::Options>,
81) -> Result<Value> {
82    ensure!(
83        params.len() == 1,
84        InvalidFuncArgsSnafu {
85            err_msg: format!(
86                "The length of the args is not correct, expect 1, have: {}",
87                params.len()
88            ),
89        }
90    );
91
92    let ValueRef::String(table_name) = params[0] else {
93        return UnsupportedInputDataTypeSnafu {
94            function,
95            datatypes: params.iter().map(|v| v.data_type()).collect::<Vec<_>>(),
96        }
97        .fail();
98    };
99
100    let (catalog_name, schema_name, table_name) = table_name_to_full_name(table_name, query_ctx)
101        .map_err(BoxedError::new)
102        .context(TableMutationSnafu)?;
103
104    let affected_rows = table_mutation_handler
105        .build_index(
106            BuildIndexTableRequest {
107                options,
108                catalog_name,
109                schema_name,
110                table_name,
111            },
112            query_ctx.clone(),
113        )
114        .await?;
115
116    Ok(Value::from(affected_rows as u64))
117}
118
119fn build_index_signature() -> Signature {
120    Signature::uniform(1, vec![ArrowDataType::Utf8], Volatility::Immutable)
121}
122
123#[cfg(test)]
124mod tests {
125    use super::*;
126    use crate::function::FunctionContext;
127
128    #[tokio::test]
129    async fn test_build_series_index_arguments() {
130        let ctx = FunctionContext::mock();
131        let handler = ctx.state.table_mutation_handler.as_ref().unwrap();
132        for params in [
133            vec![],
134            vec![ValueRef::Int32(1)],
135            vec![ValueRef::Null],
136            vec![ValueRef::String("a"), ValueRef::String("b")],
137        ] {
138            assert!(
139                build_index_impl(
140                    handler,
141                    &ctx.query_ctx,
142                    &params,
143                    "build_series_index",
144                    Some(build_index_request::Options::SeriesIndex(Default::default()))
145                )
146                .await
147                .is_err()
148            );
149        }
150        assert_eq!(
151            Value::UInt64(42),
152            build_index_impl(
153                handler,
154                &ctx.query_ctx,
155                &[ValueRef::String("greptime.public.test")],
156                "build_series_index",
157                Some(build_index_request::Options::SeriesIndex(Default::default())),
158            )
159            .await
160            .unwrap()
161        );
162    }
163}