teaql_provider_linux/
executor.rs1use 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
15pub 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 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 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 let rows = self.collect_records(entity)?;
79
80 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}