1use std::sync::{Arc, LazyLock, Weak};
36
37use common_catalog::consts::{
38 CONFIDENCE_COLUMN, DEFAULT_PRIVATE_SCHEMA_NAME, DST_ID_COLUMN, DST_TYPE_COLUMN,
39 DURATION_COUNT_COLUMN, DURATION_MAX_COLUMN, DURATION_SUM_COLUMN, EDGE_ATTRIBUTES_COLUMN,
40 ENTITY_DESCRIPTIVE_COLUMN, ENTITY_ID_ATTRS_COLUMN, ENTITY_ID_COLUMN, ENTITY_SCOPE_COLUMN,
41 ENTITY_TYPE_COLUMN, ERROR_COUNT_COLUMN, FRESH_UNTIL_COLUMN, OBSERVED_AT_COLUMN,
42 PROVENANCE_COLUMN, REL_TYPE_COLUMN, REQUEST_COUNT_COLUMN, SEMANTIC_ENTITIES_TABLE_ID,
43 SEMANTIC_ENTITIES_TABLE_NAME as SEMANTIC_ENTITIES, SEMANTIC_RELATIONSHIPS_TABLE_ID,
44 SEMANTIC_RELATIONSHIPS_TABLE_NAME as SEMANTIC_RELATIONSHIPS, SOURCE_TABLES_COLUMN,
45 SRC_ID_COLUMN, SRC_TYPE_COLUMN, UNMATCHED_COUNT_COLUMN, WINDOW_END_COLUMN, WINDOW_START_COLUMN,
46};
47use common_error::ext::BoxedError;
48use common_recordbatch::adapter::AsyncRecordBatchStreamAdapter;
49use common_recordbatch::{
50 DfRecordBatch, EmptyRecordBatchStream, RecordBatch, RecordBatchStreamWrapper,
51 SendableRecordBatchStream,
52};
53use datatypes::prelude::ConcreteDataType;
54use datatypes::schema::{ColumnSchema, Schema, SchemaRef};
55use futures::StreamExt;
56use session::context::QueryContextRef;
57use snafu::ResultExt;
58use store_api::storage::{ScanRequest, TableId};
59use table::TableRef;
60use table::metadata::TableInfo;
61
62use crate::CatalogManager;
63use crate::error::{InternalSnafu, Result};
64use crate::system_schema::{SystemSchemaProviderInner, SystemTable, SystemTableRef, utils};
65
66pub type EntityGraphProviderRef = Arc<dyn EntityGraphProvider>;
67
68pub enum DeclarationOrigin {
70 Declared,
72 Convention,
74}
75
76impl DeclarationOrigin {
77 pub fn as_str(&self) -> &'static str {
78 match self {
79 Self::Declared => "declared",
80 Self::Convention => "convention",
81 }
82 }
83}
84
85pub struct TableEntityDeclaration {
89 pub entity_type: String,
90 pub origin: DeclarationOrigin,
91 pub id_columns: Vec<String>,
92 pub id_qualifier: Option<String>,
93 pub superseded_by_columns: Vec<String>,
96 pub descriptive_columns: Vec<String>,
97 pub scope_columns: Vec<String>,
98}
99
100#[async_trait::async_trait]
113pub trait EntityGraphProvider: Send + Sync {
114 async fn scan_entities(
117 &self,
118 catalog: &str,
119 request: ScanRequest,
120 query_ctx: Option<QueryContextRef>,
121 ) -> std::result::Result<Option<SendableRecordBatchStream>, BoxedError>;
122
123 async fn scan_relationships(
125 &self,
126 catalog: &str,
127 request: ScanRequest,
128 query_ctx: Option<QueryContextRef>,
129 ) -> std::result::Result<Option<SendableRecordBatchStream>, BoxedError>;
130
131 fn table_declarations(&self, table_info: &TableInfo) -> Vec<TableEntityDeclaration>;
136}
137
138pub(crate) struct SemanticGraphTableProvider {
142 catalog_name: String,
143 catalog_manager: Weak<dyn CatalogManager>,
144 query_ctx: Option<QueryContextRef>,
147}
148
149impl SemanticGraphTableProvider {
150 pub(crate) fn new(
151 catalog_name: String,
152 catalog_manager: Weak<dyn CatalogManager>,
153 query_ctx: Option<QueryContextRef>,
154 ) -> Self {
155 Self {
156 catalog_name,
157 catalog_manager,
158 query_ctx,
159 }
160 }
161
162 pub(crate) fn table_names() -> Vec<String> {
163 vec![
164 SEMANTIC_ENTITIES.to_string(),
165 SEMANTIC_RELATIONSHIPS.to_string(),
166 ]
167 }
168
169 pub(crate) fn table_exists(name: &str) -> bool {
170 name == SEMANTIC_ENTITIES || name == SEMANTIC_RELATIONSHIPS
171 }
172
173 pub(crate) fn table(&self, name: &str) -> Option<TableRef> {
174 self.build_table(name)
175 }
176}
177
178impl SystemSchemaProviderInner for SemanticGraphTableProvider {
179 fn catalog_name(&self) -> &str {
180 &self.catalog_name
181 }
182
183 fn schema_name() -> &'static str {
184 DEFAULT_PRIVATE_SCHEMA_NAME
185 }
186
187 fn system_table(&self, name: &str) -> Option<SystemTableRef> {
188 let kind = match name {
189 SEMANTIC_ENTITIES => GraphTableKind::Entities,
190 SEMANTIC_RELATIONSHIPS => GraphTableKind::Relationships,
191 _ => return None,
192 };
193 Some(Arc::new(SemanticGraphTable::new(
194 kind,
195 self.catalog_name.clone(),
196 self.catalog_manager.clone(),
197 self.query_ctx.clone(),
198 )) as _)
199 }
200}
201
202fn ts() -> ConcreteDataType {
203 ConcreteDataType::timestamp_millisecond_datatype()
204}
205
206fn string() -> ConcreteDataType {
207 ConcreteDataType::string_datatype()
208}
209
210fn json() -> ConcreteDataType {
211 ConcreteDataType::json_datatype()
212}
213
214static ENTITIES_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
236 Arc::new(Schema::new(vec![
237 ColumnSchema::new(OBSERVED_AT_COLUMN, ts(), false).with_time_index(true),
238 ColumnSchema::new(WINDOW_START_COLUMN, ts(), true),
239 ColumnSchema::new(WINDOW_END_COLUMN, ts(), true),
240 ColumnSchema::new(FRESH_UNTIL_COLUMN, ts(), true),
241 ColumnSchema::new(ENTITY_TYPE_COLUMN, string(), false),
242 ColumnSchema::new(ENTITY_ID_COLUMN, string(), false),
243 ColumnSchema::new(ENTITY_ID_ATTRS_COLUMN, json(), true),
244 ColumnSchema::new(ENTITY_SCOPE_COLUMN, string(), true),
245 ColumnSchema::new(ENTITY_DESCRIPTIVE_COLUMN, json(), true),
246 ColumnSchema::new(SOURCE_TABLES_COLUMN, json(), true),
247 ]))
248});
249
250static RELATIONSHIPS_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
283 Arc::new(Schema::new(vec![
284 ColumnSchema::new(OBSERVED_AT_COLUMN, ts(), false).with_time_index(true),
285 ColumnSchema::new(WINDOW_START_COLUMN, ts(), true),
286 ColumnSchema::new(WINDOW_END_COLUMN, ts(), true),
287 ColumnSchema::new(FRESH_UNTIL_COLUMN, ts(), true),
288 ColumnSchema::new(SRC_TYPE_COLUMN, string(), false),
289 ColumnSchema::new(SRC_ID_COLUMN, string(), false),
290 ColumnSchema::new(DST_TYPE_COLUMN, string(), false),
291 ColumnSchema::new(DST_ID_COLUMN, string(), false),
292 ColumnSchema::new(REL_TYPE_COLUMN, string(), false),
293 ColumnSchema::new(PROVENANCE_COLUMN, string(), false),
294 ColumnSchema::new(
295 CONFIDENCE_COLUMN,
296 ConcreteDataType::float64_datatype(),
297 true,
298 ),
299 ColumnSchema::new(
300 REQUEST_COUNT_COLUMN,
301 ConcreteDataType::int64_datatype(),
302 true,
303 ),
304 ColumnSchema::new(
305 UNMATCHED_COUNT_COLUMN,
306 ConcreteDataType::int64_datatype(),
307 true,
308 ),
309 ColumnSchema::new(ERROR_COUNT_COLUMN, ConcreteDataType::int64_datatype(), true),
310 ColumnSchema::new(
311 DURATION_SUM_COLUMN,
312 ConcreteDataType::float64_datatype(),
313 true,
314 ),
315 ColumnSchema::new(
316 DURATION_COUNT_COLUMN,
317 ConcreteDataType::int64_datatype(),
318 true,
319 ),
320 ColumnSchema::new(
321 DURATION_MAX_COLUMN,
322 ConcreteDataType::float64_datatype(),
323 true,
324 ),
325 ColumnSchema::new(EDGE_ATTRIBUTES_COLUMN, json(), true),
326 ]))
327});
328
329#[derive(Clone, Copy)]
331enum GraphTableKind {
332 Entities,
333 Relationships,
334}
335
336struct SemanticGraphTable {
338 kind: GraphTableKind,
339 schema: SchemaRef,
340 catalog_name: String,
341 catalog_manager: Weak<dyn CatalogManager>,
342 query_ctx: Option<QueryContextRef>,
343}
344
345impl SemanticGraphTable {
346 fn new(
347 kind: GraphTableKind,
348 catalog_name: String,
349 catalog_manager: Weak<dyn CatalogManager>,
350 query_ctx: Option<QueryContextRef>,
351 ) -> Self {
352 let schema = match kind {
353 GraphTableKind::Entities => ENTITIES_SCHEMA.clone(),
354 GraphTableKind::Relationships => RELATIONSHIPS_SCHEMA.clone(),
355 };
356 Self {
357 kind,
358 schema,
359 catalog_name,
360 catalog_manager,
361 query_ctx,
362 }
363 }
364
365 async fn derive(
366 kind: GraphTableKind,
367 catalog: String,
368 catalog_manager: Weak<dyn CatalogManager>,
369 request: ScanRequest,
370 query_ctx: Option<QueryContextRef>,
371 ) -> Result<Option<SendableRecordBatchStream>> {
372 let provider = utils::entity_graph_provider(&catalog_manager)?;
373 let Some(provider) = provider else {
375 return Ok(None);
376 };
377 match kind {
378 GraphTableKind::Entities => provider.scan_entities(&catalog, request, query_ctx).await,
379 GraphTableKind::Relationships => {
380 provider
381 .scan_relationships(&catalog, request, query_ctx)
382 .await
383 }
384 }
385 .context(InternalSnafu)
386 }
387
388 fn align_schema(
389 stream: SendableRecordBatchStream,
390 schema: SchemaRef,
391 ) -> SendableRecordBatchStream {
392 let batch_schema = schema.clone();
393 let arrow_schema = schema.arrow_schema().clone();
394 let batches = stream.map(move |batch| {
395 let batch = batch?;
396 let batch = DfRecordBatch::try_new(
399 arrow_schema.clone(),
400 batch.into_df_record_batch().columns().to_vec(),
401 )
402 .context(common_recordbatch::error::NewDfRecordBatchSnafu)?;
403 Ok(RecordBatch::from_df_record_batch(
404 batch_schema.clone(),
405 batch,
406 ))
407 });
408 Box::pin(RecordBatchStreamWrapper::new(schema, Box::pin(batches)))
409 }
410}
411
412impl SystemTable for SemanticGraphTable {
413 fn table_id(&self) -> TableId {
414 match self.kind {
415 GraphTableKind::Entities => SEMANTIC_ENTITIES_TABLE_ID,
416 GraphTableKind::Relationships => SEMANTIC_RELATIONSHIPS_TABLE_ID,
417 }
418 }
419
420 fn table_name(&self) -> &'static str {
421 match self.kind {
422 GraphTableKind::Entities => SEMANTIC_ENTITIES,
423 GraphTableKind::Relationships => SEMANTIC_RELATIONSHIPS,
424 }
425 }
426
427 fn schema(&self) -> SchemaRef {
428 self.schema.clone()
429 }
430
431 fn to_stream(&self, request: ScanRequest) -> Result<SendableRecordBatchStream> {
432 let schema = self.schema.clone();
433 let kind = self.kind;
434 let catalog = self.catalog_name.clone();
435 let catalog_manager = self.catalog_manager.clone();
436 let query_ctx = self.query_ctx.clone();
437
438 let stream_schema = schema.clone();
439 let stream = async move {
440 let stream = Self::derive(kind, catalog, catalog_manager, request, query_ctx)
441 .await
442 .map_err(BoxedError::new)
443 .context(common_recordbatch::error::ExternalSnafu)?;
444 Ok(match stream {
445 Some(stream) => Self::align_schema(stream, stream_schema.clone()),
446 None => Box::pin(EmptyRecordBatchStream::new(stream_schema.clone())),
447 })
448 };
449
450 Ok(Box::pin(AsyncRecordBatchStreamAdapter::new(
451 schema,
452 Box::pin(stream),
453 )))
454 }
455}
456
457#[cfg(test)]
458mod tests {
459 use super::*;
460
461 #[test]
462 fn graph_tables_use_observed_at_as_time_index() {
463 for schema in [&*ENTITIES_SCHEMA, &*RELATIONSHIPS_SCHEMA] {
464 assert_eq!(
465 schema.timestamp_column().map(|column| column.name.as_str()),
466 Some("observed_at")
467 );
468 }
469 }
470
471 #[test]
472 fn relationship_schema_does_not_expose_generation_id() {
473 assert!(
474 RELATIONSHIPS_SCHEMA
475 .column_schema_by_name("generation_id")
476 .is_none()
477 );
478 }
479}