frontend/instance/
opentsdb.rs1use std::sync::Arc;
16
17use async_trait::async_trait;
18use auth::{
19 OPENTSDB_WRITE, PermissionChecker, PermissionCheckerRef, PermissionReq, PermissionTableTarget,
20 PermissionTableTargets,
21};
22use common_error::ext::BoxedError;
23use common_telemetry::tracing;
24use servers::error::{self as server_error, AuthSnafu, ExecuteGrpcQuerySnafu};
25use servers::opentsdb::codec::DataPoint;
26use servers::opentsdb::data_point_to_grpc_row_insert_requests;
27use servers::query_handler::OpentsdbProtocolHandler;
28use session::context::QueryContextRef;
29use snafu::prelude::*;
30use table::requests::{SEMANTIC_SIGNAL_TYPE, SEMANTIC_SOURCE, SIGNAL_TYPE_METRIC, SOURCE_OPENTSDB};
31
32use crate::instance::Instance;
33
34fn permission_targets(data_points: &[DataPoint], ctx: &QueryContextRef) -> PermissionTableTargets {
35 let catalog = ctx.current_catalog();
36 let schema = ctx.current_schema();
37 PermissionTableTargets::resolved(
38 data_points
39 .iter()
40 .map(|data_point| PermissionTableTarget::new(catalog, &schema, data_point.metric()))
41 .collect(),
42 )
43}
44
45#[async_trait]
46impl OpentsdbProtocolHandler for Instance {
47 async fn preflight(
48 &self,
49 data_points: &[DataPoint],
50 ctx: QueryContextRef,
51 ) -> server_error::Result<()> {
52 self.check_table_permission(
53 &ctx,
54 PermissionReq::Action(OPENTSDB_WRITE),
55 permission_targets(data_points, &ctx),
56 )
57 .context(AuthSnafu)?;
58 Ok(())
59 }
60
61 #[tracing::instrument(skip_all, fields(protocol = "opentsdb"))]
62 async fn exec(
63 &self,
64 data_points: Vec<DataPoint>,
65 ctx: QueryContextRef,
66 ) -> server_error::Result<usize> {
67 self.plugins
68 .get::<PermissionCheckerRef>()
69 .as_ref()
70 .check_permission(ctx.current_user(), PermissionReq::Action(OPENTSDB_WRITE))
71 .context(AuthSnafu)?;
72
73 let (requests, _) = data_point_to_grpc_row_insert_requests(data_points)?;
74 self.check_row_insert_permission(&requests, &ctx, PermissionReq::Action(OPENTSDB_WRITE))
75 .context(AuthSnafu)?;
76
77 let ctx = {
78 let mut c = (*ctx).clone();
79 c.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_METRIC);
80 c.set_extension(SEMANTIC_SOURCE, SOURCE_OPENTSDB);
81 Arc::new(c)
82 };
83
84 let output = self
86 .handle_row_inserts(requests, ctx, true, true)
87 .await
88 .map_err(BoxedError::new)
89 .context(ExecuteGrpcQuerySnafu)?;
90
91 Ok(match output.data {
92 common_query::OutputData::AffectedRows(rows) => rows,
93 _ => unreachable!(),
94 })
95 }
96}
97
98#[cfg(test)]
99mod tests {
100 use session::context::QueryContext;
101
102 use super::*;
103
104 #[test]
105 fn test_permission_targets_do_not_require_row_conversion() {
106 let data_points = [
107 DataPoint::new("cpu".to_string(), 0, 1.0, vec![]),
108 DataPoint::new(
109 "mem".to_string(),
110 0,
111 1.0,
112 vec![("greptime_value".to_string(), "tag".to_string())],
113 ),
114 ];
115 assert!(data_point_to_grpc_row_insert_requests(data_points.to_vec()).is_err());
116
117 assert_eq!(
118 PermissionTableTargets::Resolved(vec![
119 PermissionTableTarget::new("greptime", "public", "cpu"),
120 PermissionTableTarget::new("greptime", "public", "mem"),
121 ]),
122 permission_targets(&data_points, &QueryContext::arc())
123 );
124 }
125}