1pub 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#[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 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 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 async fn pre_write(&self, request: &RowInsertRequests, ctx: QueryContextRef) -> Result<()>;
113
114 async fn write_prepared(
116 &self,
117 request: RowInsertRequests,
118 ctx: QueryContextRef,
119 with_metric_engine: bool,
120 ) -> Result<Output>;
121
122 async fn write(
124 &self,
125 request: RowInsertRequests,
126 ctx: QueryContextRef,
127 with_metric_engine: bool,
128 ) -> Result<Output>;
129
130 async fn write_all(
132 &self,
133 requests: Vec<(QueryContextRef, RowInsertRequests)>,
134 with_metric_engine: bool,
135 ) -> Result<Vec<Result<Output>>>;
136
137 async fn read(&self, request: ReadRequest, ctx: QueryContextRef) -> Result<PromStoreResponse>;
139}
140
141#[async_trait]
142pub trait OpenTelemetryProtocolHandler: PipelineHandler {
143 async fn metrics(
145 &self,
146 request: ExportMetricsServiceRequest,
147 ctx: QueryContextRef,
148 ) -> Result<MetricsIngestOutcome>;
149
150 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#[async_trait]
180pub trait PipelineHandler {
181 async fn insert(&self, input: RowInsertRequests, ctx: QueryContextRef) -> Result<Output>;
182
183 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 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 fn build_pipeline(&self, pipeline: &str) -> Result<Pipeline>;
227
228 async fn get_pipeline_str(
230 &self,
231 name: &str,
232 version: PipelineVersion,
233 query_ctx: QueryContextRef,
234 ) -> Result<(String, TimestampNanosecond)>;
235}
236
237pub 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#[async_trait]
251pub trait LogQueryHandler {
252 async fn query(&self, query: LogQuery, ctx: QueryContextRef) -> Result<Output>;
254
255 fn catalog_manager(&self, ctx: &QueryContext) -> Result<&dyn CatalogManager>;
257}
258
259#[async_trait]
261pub trait JaegerQueryHandler {
262 async fn get_services(&self, ctx: QueryContextRef) -> Result<Output>;
264
265 async fn get_operations(
267 &self,
268 ctx: QueryContextRef,
269 service_name: &str,
270 span_kind: Option<&str>,
271 ) -> Result<Output>;
272
273 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 async fn find_traces(
288 &self,
289 ctx: QueryContextRef,
290 query_params: QueryTraceParams,
291 ) -> Result<Output>;
292}