Skip to main content

laminar_db/
show_commands.rs

1//! SHOW and DESCRIBE command builders; reopens `impl LaminarDB` to keep `db.rs` focused.
2
3use 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    /// Build a SHOW CHECKPOINT STATUS metadata result. Public so the server's
13    /// `GET /api/v1/cluster/checkpoints` endpoint can surface it.
14    ///
15    /// # Errors
16    ///
17    /// Returns [`DbError::Checkpoint`] if the metadata batch cannot be assembled.
18    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    /// Build a SHOW MATERIALIZED VIEWS metadata result.
64    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    /// Build a SHOW SOURCES metadata result with connector metadata.
95    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    /// Build a SHOW SINKS metadata result with connector metadata.
136    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    /// Build a SHOW QUERIES metadata result.
179    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    /// Build a SHOW STREAMS metadata result with SQL definitions.
201    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    /// Build a SHOW CREATE SOURCE result.
229    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    /// Build a SHOW CREATE SINK result.
253    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    /// Build a SHOW TABLES metadata result.
277    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    /// Build a DESCRIBE result.
319    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
347/// Convert an Arrow schema to a DESCRIBE result `RecordBatch`.
348pub(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}