Skip to main content

teaql_provider_linux/
executor.rs

1use std::collections::HashMap;
2use std::time::SystemTime;
3
4use teaql_core::{EntityDescriptor, Record};
5use teaql_data_service::{
6    DataServiceCapabilities, DataServiceExecutor, DataServiceOperation, ExecutionMetadata,
7    MutationExecutor, MutationRequest, MutationResult, QueryExecutor, QueryRequest, QueryResult,
8    StreamChunk, StreamQueryExecutor, Transaction, TransactionExecutor,
9};
10use teaql_runtime::{InMemoryMetadataStore, MemoryDataService};
11
12use crate::collector::{Collector, ProcessCollector, SystemInfoCollector, ThreadCollector};
13use crate::error::LinuxProviderError;
14
15/// A `DataServiceExecutor` backed by Linux /proc collectors.
16///
17/// Routes queries to the appropriate collector by entity name, collects all records,
18/// then delegates in-memory query processing (filter, sort, project, aggregate) to
19/// `MemoryDataService`.
20pub struct LinuxDataServiceExecutor {
21    collectors: HashMap<String, Box<dyn Collector>>,
22}
23
24impl Default for LinuxDataServiceExecutor {
25    fn default() -> Self {
26        Self::new()
27    }
28}
29
30impl LinuxDataServiceExecutor {
31    /// Create a new executor with the default set of collectors.
32    pub fn new() -> Self {
33        let mut collectors: HashMap<String, Box<dyn Collector>> = HashMap::new();
34        collectors.insert("SystemInfo".to_owned(), Box::new(SystemInfoCollector));
35        collectors.insert("Process".to_owned(), Box::new(ProcessCollector));
36        collectors.insert("Thread".to_owned(), Box::new(ThreadCollector));
37        Self { collectors }
38    }
39
40    /// Register an additional collector.
41    pub fn with_collector(mut self, collector: Box<dyn Collector>) -> Self {
42        self.collectors
43            .insert(collector.entity_name().to_owned(), collector);
44        self
45    }
46
47    fn collect_records(&self, entity: &str) -> Result<Vec<Record>, LinuxProviderError> {
48        let collector = self
49            .collectors
50            .get(entity)
51            .ok_or_else(|| LinuxProviderError::UnknownEntity(entity.to_owned()))?;
52        collector.collect_all()
53    }
54}
55
56impl DataServiceExecutor for LinuxDataServiceExecutor {
57    type Error = LinuxProviderError;
58
59    fn capabilities(&self) -> DataServiceCapabilities {
60        DataServiceCapabilities {
61            query: true,
62            mutation: false,
63            transaction: false,
64            schema: false,
65            id_generation: false,
66            batch_mutation: false,
67            returning: false,
68        }
69    }
70}
71
72impl QueryExecutor for LinuxDataServiceExecutor {
73    async fn query(&self, request: QueryRequest) -> Result<QueryResult, Self::Error> {
74        let started_at = SystemTime::now();
75        let entity = &request.query.entity;
76
77        // Collect raw records from the matching collector.
78        let rows = self.collect_records(entity)?;
79
80        // Build a MemoryDataService with the entity registered so fetch_all succeeds.
81        let metadata = InMemoryMetadataStore::new().with_entity(EntityDescriptor::new(entity));
82        let mut mem = MemoryDataService::new(metadata);
83        mem.seed(entity.clone(), rows);
84
85        let result_rows = mem
86            .fetch_all(&request.query)
87            .map_err(|e| LinuxProviderError::ProcFs(e.to_string()))?;
88
89        let ended_at = SystemTime::now();
90        let count = result_rows.len();
91
92        Ok(QueryResult {
93            rows: result_rows,
94            metadata: ExecutionMetadata {
95                backend: "linux-proc".to_owned(),
96                operation: DataServiceOperation::Query,
97                started_at,
98                ended_at,
99                affected_rows: None,
100                result_count: Some(count),
101                trace_chain: request.trace_chain,
102                comment: request.comment,
103                backend_request_id: None,
104                debug_query: None,
105            },
106        })
107    }
108}
109
110impl MutationExecutor for LinuxDataServiceExecutor {
111    async fn mutate(&self, _request: MutationRequest) -> Result<MutationResult, Self::Error> {
112        Err(LinuxProviderError::ProcFs(
113            "Linux provider is read-only".to_owned(),
114        ))
115    }
116}
117
118impl Transaction for LinuxDataServiceExecutor {
119    type Error = LinuxProviderError;
120
121    async fn commit(self) -> Result<(), Self::Error> {
122        Ok(())
123    }
124    async fn rollback(self) -> Result<(), Self::Error> {
125        Ok(())
126    }
127}
128
129impl TransactionExecutor for LinuxDataServiceExecutor {
130    type Tx<'a> = LinuxDataServiceExecutor;
131
132    async fn begin(&self) -> Result<Self::Tx<'_>, Self::Error> {
133        Err(LinuxProviderError::ProcFs(
134            "Linux provider does not support transactions".to_owned(),
135        ))
136    }
137}
138
139impl StreamQueryExecutor for LinuxDataServiceExecutor {
140    async fn query_stream(
141        &self,
142        _request: QueryRequest,
143        _chunk_size: usize,
144    ) -> Result<Vec<StreamChunk>, Self::Error> {
145        Err(LinuxProviderError::ProcFs(
146            "Linux provider does not support streaming".to_owned(),
147        ))
148    }
149}