1use std::sync::Arc;
4
5use arrow::array::{BooleanArray, RecordBatch, StringArray, UInt64Array};
6use arrow::datatypes::{DataType, Field, Schema};
7
8use crate::db::LaminarDB;
9use crate::error::DbError;
10
11impl LaminarDB {
12 pub async fn build_show_checkpoint_status(&self) -> Result<RecordBatch, DbError> {
19 let coordinator = self.coordinator.lock().await;
20 let (cp_id, epoch, ts_ms, sources, sinks, completed_this_runtime) =
21 if let Some(coordinator) = coordinator.as_ref() {
22 let completed = coordinator.stats().completed;
23 coordinator.last_committed_manifest().map_or_else(
24 || (0, 0, 0, String::new(), String::new(), completed),
25 |manifest| {
26 (
27 manifest.checkpoint_id,
28 manifest.epoch,
29 manifest.timestamp_ms,
30 manifest.source_names.join(", "),
31 manifest.sink_names.join(", "),
32 completed,
33 )
34 },
35 )
36 } else {
37 (0, 0, 0, String::new(), String::new(), 0)
38 };
39
40 let schema = Arc::new(Schema::new(vec![
41 Field::new("checkpoint_id", DataType::UInt64, false),
42 Field::new("epoch", DataType::UInt64, false),
43 Field::new("timestamp_ms", DataType::UInt64, false),
44 Field::new("sources", DataType::Utf8, false),
45 Field::new("sinks", DataType::Utf8, false),
46 Field::new("completed_this_runtime", DataType::UInt64, false),
47 ]));
48
49 RecordBatch::try_new(
50 schema,
51 vec![
52 Arc::new(UInt64Array::from(vec![cp_id])),
53 Arc::new(UInt64Array::from(vec![epoch])),
54 Arc::new(UInt64Array::from(vec![ts_ms])),
55 Arc::new(StringArray::from(vec![sources.as_str()])),
56 Arc::new(StringArray::from(vec![sinks.as_str()])),
57 Arc::new(UInt64Array::from(vec![completed_this_runtime])),
58 ],
59 )
60 .map_err(|e| DbError::Checkpoint(format!("failed to build checkpoint status batch: {e}")))
61 }
62
63 pub(crate) fn build_show_materialized_views(&self) -> RecordBatch {
65 let registry = self.mv_registry.lock();
66 let mut names = Vec::new();
67 let mut sqls = Vec::new();
68 let mut states = Vec::new();
69 for view in registry.views() {
70 let info = crate::handle::MaterializedViewInfo::from(view);
71 names.push(info.name);
72 sqls.push(info.sql);
73 states.push(info.state);
74 }
75 let names_ref: Vec<&str> = names.iter().map(String::as_str).collect();
76 let sqls_ref: Vec<&str> = sqls.iter().map(String::as_str).collect();
77 let states_ref: Vec<&str> = states.iter().map(String::as_str).collect();
78 let schema = Arc::new(Schema::new(vec![
79 Field::new("view_name", DataType::Utf8, false),
80 Field::new("sql", DataType::Utf8, false),
81 Field::new("state", DataType::Utf8, false),
82 ]));
83 RecordBatch::try_new(
84 schema,
85 vec![
86 Arc::new(StringArray::from(names_ref)),
87 Arc::new(StringArray::from(sqls_ref)),
88 Arc::new(StringArray::from(states_ref)),
89 ],
90 )
91 .expect("show materialized views: schema matches columns")
92 }
93
94 pub(crate) fn build_show_sources(&self) -> RecordBatch {
96 let sources = self.sources();
97 let mgr = self.connector_manager.lock();
98 let regs = mgr.sources();
99
100 let mut names = Vec::with_capacity(sources.len());
101 let mut connectors: Vec<Option<&str>> = Vec::with_capacity(sources.len());
102 let mut formats: Vec<Option<&str>> = Vec::with_capacity(sources.len());
103 let mut watermarks: Vec<Option<&str>> = Vec::with_capacity(sources.len());
104
105 for s in &sources {
106 names.push(s.name.as_str());
107 if let Some(reg) = regs.get(&s.name) {
108 connectors.push(reg.connector_type.as_deref());
109 formats.push(reg.format.as_deref());
110 } else {
111 connectors.push(None);
112 formats.push(None);
113 }
114 watermarks.push(s.watermark_column.as_deref());
115 }
116
117 let schema = Arc::new(Schema::new(vec![
118 Field::new("source_name", DataType::Utf8, false),
119 Field::new("connector", DataType::Utf8, true),
120 Field::new("format", DataType::Utf8, true),
121 Field::new("watermark_column", DataType::Utf8, true),
122 ]));
123 RecordBatch::try_new(
124 schema,
125 vec![
126 Arc::new(StringArray::from(names)),
127 Arc::new(StringArray::from(connectors)),
128 Arc::new(StringArray::from(formats)),
129 Arc::new(StringArray::from(watermarks)),
130 ],
131 )
132 .expect("show sources: schema matches columns")
133 }
134
135 pub(crate) fn build_show_sinks(&self) -> RecordBatch {
137 let sinks = self.sinks();
138 let mgr = self.connector_manager.lock();
139 let regs = mgr.sinks();
140
141 let mut names = Vec::with_capacity(sinks.len());
142 let mut inputs: Vec<Option<String>> = Vec::with_capacity(sinks.len());
143 let mut connectors: Vec<Option<&str>> = Vec::with_capacity(sinks.len());
144 let mut formats: Vec<Option<&str>> = Vec::with_capacity(sinks.len());
145
146 for s in &sinks {
147 names.push(s.name.as_str());
148 let catalog_input = self.catalog.get_sink_input(&s.name);
149 if let Some(reg) = regs.get(&s.name) {
150 inputs.push(Some(reg.input.clone()));
151 connectors.push(reg.connector_type.as_deref());
152 formats.push(reg.format.as_deref());
153 } else {
154 inputs.push(catalog_input);
155 connectors.push(None);
156 formats.push(None);
157 }
158 }
159
160 let schema = Arc::new(Schema::new(vec![
161 Field::new("sink_name", DataType::Utf8, false),
162 Field::new("input", DataType::Utf8, true),
163 Field::new("connector", DataType::Utf8, true),
164 Field::new("format", DataType::Utf8, true),
165 ]));
166 RecordBatch::try_new(
167 schema,
168 vec![
169 Arc::new(StringArray::from(names)),
170 Arc::new(StringArray::from(inputs)),
171 Arc::new(StringArray::from(connectors)),
172 Arc::new(StringArray::from(formats)),
173 ],
174 )
175 .expect("show sinks: schema matches columns")
176 }
177
178 pub(crate) fn build_show_queries(&self) -> RecordBatch {
180 let queries = self.queries();
181 let ids: Vec<u64> = queries.iter().map(|q| q.id).collect();
182 let sqls: Vec<&str> = queries.iter().map(|q| q.sql.as_str()).collect();
183 let actives: Vec<bool> = queries.iter().map(|q| q.active).collect();
184 let schema = Arc::new(Schema::new(vec![
185 Field::new("query_id", DataType::UInt64, false),
186 Field::new("sql", DataType::Utf8, false),
187 Field::new("active", DataType::Boolean, false),
188 ]));
189 RecordBatch::try_new(
190 schema,
191 vec![
192 Arc::new(UInt64Array::from(ids)),
193 Arc::new(StringArray::from(sqls)),
194 Arc::new(BooleanArray::from(actives)),
195 ],
196 )
197 .expect("show queries: schema matches columns")
198 }
199
200 pub(crate) fn build_show_streams(&self) -> RecordBatch {
202 let streams = self.catalog.list_streams();
203 let mgr = self.connector_manager.lock();
204 let regs = mgr.streams();
205
206 let mut names = Vec::with_capacity(streams.len());
207 let mut sqls: Vec<Option<&str>> = Vec::with_capacity(streams.len());
208
209 for name in &streams {
210 names.push(name.as_str());
211 sqls.push(regs.get(name.as_str()).map(|r| r.query_sql.as_str()));
212 }
213
214 let schema = Arc::new(Schema::new(vec![
215 Field::new("stream_name", DataType::Utf8, false),
216 Field::new("sql", DataType::Utf8, true),
217 ]));
218 RecordBatch::try_new(
219 schema,
220 vec![
221 Arc::new(StringArray::from(names)),
222 Arc::new(StringArray::from(sqls)),
223 ],
224 )
225 .expect("show streams: schema matches columns")
226 }
227
228 pub(crate) fn build_show_create_source(&self, name: &str) -> Result<RecordBatch, DbError> {
230 if self.catalog.get_source(name).is_none() {
231 return Err(DbError::SourceNotFound(name.to_string()));
232 }
233 let mgr = self.connector_manager.lock();
234 let ddl = mgr
235 .get_ddl(name)
236 .ok_or_else(|| DbError::InvalidOperation(format!("No stored DDL for source '{name}'")))?
237 .to_string();
238 drop(mgr);
239
240 let schema = Arc::new(Schema::new(vec![Field::new(
241 "create_statement",
242 DataType::Utf8,
243 false,
244 )]));
245 Ok(RecordBatch::try_new(
246 schema,
247 vec![Arc::new(StringArray::from(vec![ddl.as_str()]))],
248 )
249 .expect("show create source: schema matches columns"))
250 }
251
252 pub(crate) fn build_show_create_sink(&self, name: &str) -> Result<RecordBatch, DbError> {
254 if self.catalog.get_sink_input(name).is_none() {
255 return Err(DbError::SinkNotFound(name.to_string()));
256 }
257 let mgr = self.connector_manager.lock();
258 let ddl = mgr
259 .get_ddl(name)
260 .ok_or_else(|| DbError::InvalidOperation(format!("No stored DDL for sink '{name}'")))?
261 .to_string();
262 drop(mgr);
263
264 let schema = Arc::new(Schema::new(vec![Field::new(
265 "create_statement",
266 DataType::Utf8,
267 false,
268 )]));
269 Ok(RecordBatch::try_new(
270 schema,
271 vec![Arc::new(StringArray::from(vec![ddl.as_str()]))],
272 )
273 .expect("show create sink: schema matches columns"))
274 }
275
276 pub(crate) fn build_show_tables(&self) -> RecordBatch {
278 let ts = self.table_store.read();
279 let mut names = Vec::new();
280 let mut pks = Vec::new();
281 let mut row_counts = Vec::new();
282 let mut connectors = Vec::new();
283
284 for name in ts.table_names() {
285 let pk = ts.primary_key(&name).unwrap_or("").to_string();
286 let count = ts.table_row_count(&name) as u64;
287 let conn = ts.connector(&name).unwrap_or("").to_string();
288
289 names.push(name);
290 pks.push(pk);
291 row_counts.push(count);
292 connectors.push(conn);
293 }
294
295 let names_ref: Vec<&str> = names.iter().map(String::as_str).collect();
296 let pks_ref: Vec<&str> = pks.iter().map(String::as_str).collect();
297 let connectors_ref: Vec<&str> = connectors.iter().map(String::as_str).collect();
298
299 let schema = Arc::new(Schema::new(vec![
300 Field::new("name", DataType::Utf8, false),
301 Field::new("primary_key", DataType::Utf8, false),
302 Field::new("row_count", DataType::UInt64, false),
303 Field::new("connector", DataType::Utf8, false),
304 ]));
305
306 RecordBatch::try_new(
307 schema,
308 vec![
309 Arc::new(StringArray::from(names_ref)),
310 Arc::new(StringArray::from(pks_ref)),
311 Arc::new(UInt64Array::from(row_counts)),
312 Arc::new(StringArray::from(connectors_ref)),
313 ],
314 )
315 .expect("show tables: schema matches columns")
316 }
317
318 pub(crate) fn build_describe(&self, name: &str) -> Result<RecordBatch, DbError> {
320 let schema = if let Some(s) = self.catalog.describe_source(name) {
321 s
322 } else if let Some(s) = self.table_store.read().table_schema(name) {
323 s
324 } else if let Some(s) = self
325 .mv_registry
326 .lock()
327 .get(name)
328 .map(|mv| mv.schema.clone())
329 {
330 s
331 } else if self.catalog.list_sinks().contains(&name.to_string()) {
332 return Err(DbError::InvalidOperation(
333 "DESCRIBE is not supported for sinks. Use SHOW SINKS for details.".to_string(),
334 ));
335 } else if self.catalog.list_streams().contains(&name.to_string()) {
336 return Err(DbError::InvalidOperation(format!(
337 "Stream '{name}' exists but schema is only available after pipeline start"
338 )));
339 } else {
340 return Err(DbError::TableNotFound(name.to_string()));
341 };
342
343 schema_to_describe_batch(&schema)
344 }
345}
346
347pub(crate) fn schema_to_describe_batch(schema: &Schema) -> Result<RecordBatch, DbError> {
349 let col_names: Vec<String> = schema.fields().iter().map(|f| f.name().clone()).collect();
350 let col_types: Vec<String> = schema
351 .fields()
352 .iter()
353 .map(|f| format!("{}", f.data_type()))
354 .collect();
355 let col_nullable: Vec<bool> = schema.fields().iter().map(|f| f.is_nullable()).collect();
356
357 let names_ref: Vec<&str> = col_names.iter().map(String::as_str).collect();
358 let types_ref: Vec<&str> = col_types.iter().map(String::as_str).collect();
359
360 let result_schema = Arc::new(Schema::new(vec![
361 Field::new("column_name", DataType::Utf8, false),
362 Field::new("data_type", DataType::Utf8, false),
363 Field::new("nullable", DataType::Boolean, false),
364 ]));
365
366 RecordBatch::try_new(
367 result_schema,
368 vec![
369 Arc::new(StringArray::from(names_ref)),
370 Arc::new(StringArray::from(types_ref)),
371 Arc::new(BooleanArray::from(col_nullable)),
372 ],
373 )
374 .map_err(|e| DbError::InvalidOperation(format!("describe metadata: {e}")))
375}