1use std::collections::HashMap;
16use std::sync::{Arc, RwLock};
17use std::time::{Duration, Instant};
18
19use api::v1::CreateTableExpr;
20use catalog::{CatalogManagerRef, RegisterSystemTableRequest};
21use common_catalog::consts::{DEFAULT_PRIVATE_SCHEMA_NAME, default_engine};
22use common_frontend::metrics;
23use common_meta::rpc::ddl::TriggerReason;
24use common_telemetry::info;
25use common_time::FOREVER;
26use datatypes::timestamp::TimestampNanosecond;
27use futures::FutureExt;
28use operator::insert::InserterRef;
29use operator::statement::StatementExecutorRef;
30use query::QueryEngineRef;
31use session::context::QueryContextRef;
32use snafu::{OptionExt, ResultExt};
33use table::TableRef;
34use table::requests::TTL_KEY;
35
36use crate::Pipeline;
37use crate::error::{CatalogSnafu, CreateTableSnafu, PipelineTableNotFoundSnafu, Result};
38use crate::manager::{PipelineInfo, PipelineTableRef, PipelineVersion};
39use crate::metrics::{
40 METRIC_PIPELINE_CREATE_HISTOGRAM, METRIC_PIPELINE_DELETE_HISTOGRAM,
41 METRIC_PIPELINE_RETRIEVE_HISTOGRAM,
42};
43use crate::options::PipelineOptions;
44use crate::table::{PIPELINE_TABLE_NAME, PipelineTable};
45
46pub struct PipelineOperator {
55 inserter: InserterRef,
56 statement_executor: StatementExecutorRef,
57 catalog_manager: CatalogManagerRef,
58 query_engine: QueryEngineRef,
59 tables: RwLock<HashMap<String, PipelineTableRef>>,
60 cache_ttl: Duration,
61}
62
63impl PipelineOperator {
64 fn create_table_request(&self, catalog: &str) -> RegisterSystemTableRequest {
66 let (time_index, primary_keys, column_defs) = PipelineTable::build_pipeline_schema();
67
68 let mut table_options = HashMap::new();
69 table_options.insert(TTL_KEY.to_string(), FOREVER.to_string());
70
71 let create_table_expr = CreateTableExpr {
72 catalog_name: catalog.to_string(),
73 schema_name: DEFAULT_PRIVATE_SCHEMA_NAME.to_string(),
74 table_name: PIPELINE_TABLE_NAME.to_string(),
75 desc: "GreptimeDB pipeline table for Log".to_string(),
76 column_defs,
77 time_index,
78 primary_keys,
79 create_if_not_exists: true,
80 table_options,
81 table_id: None, engine: default_engine().to_string(),
83 };
84
85 RegisterSystemTableRequest {
86 create_table_expr,
87 open_hook: None,
88 }
89 }
90
91 fn add_pipeline_table_to_cache(&self, catalog: &str, table: TableRef) {
92 let mut tables = self.tables.write().unwrap();
93 if tables.contains_key(catalog) {
94 return;
95 }
96 tables.insert(
97 catalog.to_string(),
98 Arc::new(PipelineTable::new(
99 self.inserter.clone(),
100 self.statement_executor.clone(),
101 table,
102 self.query_engine.clone(),
103 self.cache_ttl,
104 )),
105 );
106 }
107
108 async fn create_pipeline_table_if_not_exists(&self, ctx: QueryContextRef) -> Result<()> {
109 let catalog = ctx.current_catalog();
110
111 if self.get_pipeline_table_from_cache(catalog).is_some() {
113 return Ok(());
114 }
115
116 let RegisterSystemTableRequest {
117 create_table_expr: mut expr,
118 open_hook: _,
119 } = self.create_table_request(catalog);
120
121 if let Some(table) = self
123 .catalog_manager
124 .table(
125 &expr.catalog_name,
126 &expr.schema_name,
127 &expr.table_name,
128 Some(&ctx),
129 )
130 .await
131 .context(CatalogSnafu)?
132 {
133 self.add_pipeline_table_to_cache(catalog, table);
134 return Ok(());
135 }
136
137 self.statement_executor
139 .create_table_inner(&mut expr, None, ctx.clone(), TriggerReason::AutoCreate)
140 .await
141 .context(CreateTableSnafu)?;
142
143 let schema = &expr.schema_name;
144 let table_name = &expr.table_name;
145
146 let table = self
148 .catalog_manager
149 .table(catalog, schema, table_name, Some(&ctx))
150 .await
151 .context(CatalogSnafu)?
152 .context(PipelineTableNotFoundSnafu)?;
153
154 info!(
155 "Created pipelines table {} with table id {}.",
156 table.table_info().full_table_name(),
157 table.table_info().table_id()
158 );
159
160 self.add_pipeline_table_to_cache(catalog, table);
162
163 Ok(())
164 }
165
166 pub fn get_pipeline_table_from_cache(&self, catalog: &str) -> Option<PipelineTableRef> {
168 let table = self.tables.read().unwrap().get(catalog).cloned();
169 metrics::record_cache_lookup("pipeline_table", table.is_some());
170 table
171 }
172}
173
174impl PipelineOperator {
175 pub fn new(
177 inserter: InserterRef,
178 statement_executor: StatementExecutorRef,
179 catalog_manager: CatalogManagerRef,
180 query_engine: QueryEngineRef,
181 options: &PipelineOptions,
182 ) -> Self {
183 Self {
184 inserter,
185 statement_executor,
186 catalog_manager,
187 tables: RwLock::new(HashMap::new()),
188 query_engine,
189 cache_ttl: options.cache_ttl,
190 }
191 }
192
193 pub async fn get_pipeline(
195 &self,
196 query_ctx: QueryContextRef,
197 name: &str,
198 version: PipelineVersion,
199 ) -> Result<Arc<Pipeline>> {
200 let schema = query_ctx.current_schema();
201 self.create_pipeline_table_if_not_exists(query_ctx.clone())
202 .await?;
203
204 let timer = Instant::now();
205 self.get_pipeline_table_from_cache(query_ctx.current_catalog())
206 .context(PipelineTableNotFoundSnafu)?
207 .get_pipeline(&schema, name, version)
208 .inspect(|re| {
209 METRIC_PIPELINE_RETRIEVE_HISTOGRAM
210 .with_label_values(&[&re.is_ok().to_string()])
211 .observe(timer.elapsed().as_secs_f64())
212 })
213 .await
214 }
215
216 pub async fn get_pipeline_str(
218 &self,
219 name: &str,
220 version: PipelineVersion,
221 query_ctx: QueryContextRef,
222 ) -> Result<(String, TimestampNanosecond)> {
223 let schema = query_ctx.current_schema();
224 self.create_pipeline_table_if_not_exists(query_ctx.clone())
225 .await?;
226
227 let timer = Instant::now();
228 self.get_pipeline_table_from_cache(query_ctx.current_catalog())
229 .context(PipelineTableNotFoundSnafu)?
230 .get_pipeline_str(&schema, name, version)
231 .inspect(|re| {
232 METRIC_PIPELINE_RETRIEVE_HISTOGRAM
233 .with_label_values(&[&re.is_ok().to_string()])
234 .observe(timer.elapsed().as_secs_f64())
235 })
236 .await
237 .map(|p| (p.content, p.version))
238 }
239
240 pub async fn insert_pipeline(
242 &self,
243 name: &str,
244 content_type: &str,
245 pipeline: &str,
246 query_ctx: QueryContextRef,
247 ) -> Result<PipelineInfo> {
248 self.create_pipeline_table_if_not_exists(query_ctx.clone())
249 .await?;
250
251 let timer = Instant::now();
252 self.get_pipeline_table_from_cache(query_ctx.current_catalog())
253 .context(PipelineTableNotFoundSnafu)?
254 .insert_and_compile(name, content_type, pipeline)
255 .inspect(|re| {
256 METRIC_PIPELINE_CREATE_HISTOGRAM
257 .with_label_values(&[&re.is_ok().to_string()])
258 .observe(timer.elapsed().as_secs_f64())
259 })
260 .await
261 }
262
263 pub async fn delete_pipeline(
265 &self,
266 name: &str,
267 version: PipelineVersion,
268 query_ctx: QueryContextRef,
269 ) -> Result<Option<()>> {
270 self.create_pipeline_table_if_not_exists(query_ctx.clone())
272 .await?;
273
274 let timer = Instant::now();
275 self.get_pipeline_table_from_cache(query_ctx.current_catalog())
276 .context(PipelineTableNotFoundSnafu)?
277 .delete_pipeline(name, version)
278 .inspect(|re| {
279 METRIC_PIPELINE_DELETE_HISTOGRAM
280 .with_label_values(&[&re.is_ok().to_string()])
281 .observe(timer.elapsed().as_secs_f64())
282 })
283 .await
284 }
285
286 pub fn build_pipeline(pipeline: &str) -> Result<Pipeline> {
288 PipelineTable::compile_pipeline(pipeline)
289 }
290}