Skip to main content

servers/
query_handler.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
15//! All query handler traits for various request protocols, like SQL or GRPC.
16//!
17//! Instance that wishes to support certain request protocol, just implement the corresponding
18//! trait, the Server will handle codec for you.
19//!
20//! Note:
21//! Query handlers are not confined to only handle read requests, they are expecting to handle
22//! write requests too. So the "query" here not might seem ambiguity. However, "query" has been
23//! used as some kind of "convention", it's the "Q" in "SQL". So we might better stick to the
24//! word "query".
25
26pub mod grpc;
27pub mod sql;
28
29use std::collections::HashMap;
30use std::sync::Arc;
31
32use api::prom_store::remote::ReadRequest;
33use api::v1::RowInsertRequests;
34use async_trait::async_trait;
35use catalog::CatalogManager;
36use common_query::Output;
37use datatypes::timestamp::TimestampNanosecond;
38use headers::HeaderValue;
39use log_query::LogQuery;
40use opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest;
41use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
42use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
43use pipeline::{GreptimePipelineParams, Pipeline, PipelineInfo, PipelineVersion, PipelineWay};
44use serde_json::Value;
45use session::context::{QueryContext, QueryContextRef};
46
47#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
48pub struct DashboardDefinition {
49    pub name: String,
50    pub definition: String,
51}
52
53use crate::error::Result;
54use crate::http::jaeger::QueryTraceParams;
55use crate::influxdb::InfluxdbRequest;
56use crate::opentsdb::codec::DataPoint;
57pub type OpentsdbProtocolHandlerRef = Arc<dyn OpentsdbProtocolHandler + Send + Sync>;
58pub type InfluxdbLineProtocolHandlerRef = Arc<dyn InfluxdbLineProtocolHandler + Send + Sync>;
59pub type PromStoreProtocolHandlerRef = Arc<dyn PromStoreProtocolHandler + Send + Sync>;
60pub type OpenTelemetryProtocolHandlerRef = Arc<dyn OpenTelemetryProtocolHandler + Send + Sync>;
61pub type PipelineHandlerRef = Arc<dyn PipelineHandler + Send + Sync>;
62pub type LogQueryHandlerRef = Arc<dyn LogQueryHandler + Send + Sync>;
63pub type JaegerQueryHandlerRef = Arc<dyn JaegerQueryHandler + Send + Sync>;
64
65#[derive(Debug, Default, Clone)]
66pub struct TraceIngestOutcome {
67    pub write_cost: usize,
68    pub accepted_spans: usize,
69    pub rejected_spans: usize,
70    pub error_message: Option<String>,
71}
72
73/// Result of ingesting one OTLP metrics request or Arrow batch.
74#[derive(Debug, Default, Clone)]
75pub struct MetricsIngestOutcome {
76    pub write_cost: usize,
77    pub accepted_data_points: i64,
78    pub rejected_data_points: i64,
79    pub error_message: Option<String>,
80}
81
82#[async_trait]
83pub trait InfluxdbLineProtocolHandler {
84    async fn exec(&self, request: InfluxdbRequest, ctx: QueryContextRef) -> Result<Output>;
85}
86
87#[async_trait]
88pub trait OpentsdbProtocolHandler {
89    /// Checks all points in one external request before per-point debug execution.
90    async fn preflight(&self, data_points: &[DataPoint], ctx: QueryContextRef) -> Result<()>;
91
92    async fn exec(&self, data_points: Vec<DataPoint>, ctx: QueryContextRef) -> Result<usize>;
93
94    /// Executes an ordinary HTTP put with optional batching. Debug callers
95    /// retain [`Self::exec`]; the frontend clears HTTP batching selection
96    /// there so diagnostic requests preserve direct, per-point error attribution.
97    async fn exec_batch(&self, data_points: Vec<DataPoint>, ctx: QueryContextRef) -> Result<usize> {
98        self.exec(data_points, ctx).await
99    }
100}
101
102pub struct PromStoreResponse {
103    pub content_type: HeaderValue,
104    pub content_encoding: HeaderValue,
105    pub resp_metrics: HashMap<String, Value>,
106    pub body: Vec<u8>,
107}
108
109#[async_trait]
110pub trait PromStoreProtocolHandler {
111    /// Runs pre-write checks/hooks for prometheus remote write requests.
112    async fn pre_write(&self, request: &RowInsertRequests, ctx: QueryContextRef) -> Result<()>;
113
114    /// Writes one batch after [`Self::pre_write`] has succeeded for the entire request.
115    async fn write_prepared(
116        &self,
117        request: RowInsertRequests,
118        ctx: QueryContextRef,
119        with_metric_engine: bool,
120    ) -> Result<Output>;
121
122    /// Handling prometheus remote write requests
123    async fn write(
124        &self,
125        request: RowInsertRequests,
126        ctx: QueryContextRef,
127        with_metric_engine: bool,
128    ) -> Result<Output>;
129
130    /// Checks every batch before writing any of them.
131    async fn write_all(
132        &self,
133        requests: Vec<(QueryContextRef, RowInsertRequests)>,
134        with_metric_engine: bool,
135    ) -> Result<Vec<Result<Output>>>;
136
137    /// Handling prometheus remote read requests
138    async fn read(&self, request: ReadRequest, ctx: QueryContextRef) -> Result<PromStoreResponse>;
139}
140
141#[async_trait]
142pub trait OpenTelemetryProtocolHandler: PipelineHandler {
143    /// Handling opentelemetry metrics request
144    async fn metrics(
145        &self,
146        request: ExportMetricsServiceRequest,
147        ctx: QueryContextRef,
148    ) -> Result<MetricsIngestOutcome>;
149
150    /// Handling opentelemetry traces request
151    async fn traces(
152        &self,
153        pipeline_handler: PipelineHandlerRef,
154        request: ExportTraceServiceRequest,
155        pipeline: PipelineWay,
156        pipeline_params: GreptimePipelineParams,
157        table_name: String,
158        ctx: QueryContextRef,
159    ) -> Result<TraceIngestOutcome>;
160
161    async fn logs(
162        &self,
163        pipeline_handler: PipelineHandlerRef,
164        request: ExportLogsServiceRequest,
165        pipeline: PipelineWay,
166        pipeline_params: GreptimePipelineParams,
167        table_name: String,
168        ctx: QueryContextRef,
169    ) -> Result<Vec<Output>>;
170}
171
172/// PipelineHandler is responsible for handling pipeline related requests.
173///
174/// The "Pipeline" is a series of transformations that can be applied to unstructured
175/// data like logs. This handler is responsible to manage pipelines and accept data for
176/// processing.
177///
178/// The pipeline is stored in the database and can be retrieved by its name.
179#[async_trait]
180pub trait PipelineHandler {
181    async fn insert(&self, input: RowInsertRequests, ctx: QueryContextRef) -> Result<Output>;
182
183    /// Checks every batch before inserting any of them.
184    async fn insert_all(
185        &self,
186        inputs: Vec<(QueryContextRef, RowInsertRequests)>,
187    ) -> Result<Vec<Result<Output>>>;
188
189    fn check_pipeline_query_permission(&self, query_ctx: &QueryContextRef) -> Result<()>;
190
191    /// Loads a compiled pipeline for execution.
192    ///
193    /// This intentionally does not check pipeline-query permission: users with
194    /// write-only permission can ingest through an existing pipeline. Inspection
195    /// and preview callers must check query permission first; ingestion enforces
196    /// write and table-target permissions separately.
197    async fn get_pipeline(
198        &self,
199        name: &str,
200        version: PipelineVersion,
201        query_ctx: QueryContextRef,
202    ) -> Result<Arc<Pipeline>>;
203
204    async fn insert_pipeline(
205        &self,
206        name: &str,
207        content_type: &str,
208        pipeline: &str,
209        query_ctx: QueryContextRef,
210    ) -> Result<PipelineInfo>;
211
212    async fn delete_pipeline(
213        &self,
214        name: &str,
215        version: PipelineVersion,
216        query_ctx: QueryContextRef,
217    ) -> Result<Option<()>>;
218
219    async fn get_table(
220        &self,
221        table: &str,
222        query_ctx: &QueryContext,
223    ) -> std::result::Result<Option<Arc<table::Table>>, catalog::error::Error>;
224
225    //// Build a pipeline from a string.
226    fn build_pipeline(&self, pipeline: &str) -> Result<Pipeline>;
227
228    /// Get a original pipeline by name.
229    async fn get_pipeline_str(
230        &self,
231        name: &str,
232        version: PipelineVersion,
233        query_ctx: QueryContextRef,
234    ) -> Result<(String, TimestampNanosecond)>;
235}
236
237/// Handling dashboard as code CRUD
238pub type DashboardHandlerRef = Arc<dyn DashboardHandler + Send + Sync>;
239
240#[async_trait]
241pub trait DashboardHandler {
242    async fn save(&self, name: &str, definition: &str, ctx: QueryContextRef) -> Result<()>;
243
244    async fn list(&self, ctx: QueryContextRef) -> Result<Vec<DashboardDefinition>>;
245
246    async fn delete(&self, name: &str, ctx: QueryContextRef) -> Result<()>;
247}
248
249/// Handle log query requests.
250#[async_trait]
251pub trait LogQueryHandler {
252    /// Execute a log query.
253    async fn query(&self, query: LogQuery, ctx: QueryContextRef) -> Result<Output>;
254
255    /// Get catalog manager.
256    fn catalog_manager(&self, ctx: &QueryContext) -> Result<&dyn CatalogManager>;
257}
258
259/// Handle Jaeger query requests.
260#[async_trait]
261pub trait JaegerQueryHandler {
262    /// Get trace services. It's used for `/api/services` API.
263    async fn get_services(&self, ctx: QueryContextRef) -> Result<Output>;
264
265    /// Get Jaeger operations. It's used for `/api/operations` and `/api/services/{service_name}/operations` API.
266    async fn get_operations(
267        &self,
268        ctx: QueryContextRef,
269        service_name: &str,
270        span_kind: Option<&str>,
271    ) -> Result<Output>;
272
273    /// Retrieves a trace by its unique identifier.
274    ///
275    /// This method is used to handle requests to the `/api/traces/{trace_id}` endpoint.
276    /// It accepts optional `start_time` and `end_time` parameters in nanoseconds to filter the trace data within a specific time range.
277    async fn get_trace(
278        &self,
279        ctx: QueryContextRef,
280        trace_id: &str,
281        start_time: Option<i64>,
282        end_time: Option<i64>,
283        limit: Option<usize>,
284    ) -> Result<Output>;
285
286    /// Find traces by query params. It's used for `/api/traces` API.
287    async fn find_traces(
288        &self,
289        ctx: QueryContextRef,
290        query_params: QueryTraceParams,
291    ) -> Result<Output>;
292}