Skip to main content

teaql_runtime/data_service/
base.rs

1use teaql_core::{
2    BatchInsertCommand, BatchUpdateCommand, DeleteCommand, Entity, InsertCommand, Record,
3    RecoverCommand, SelectQuery, SmartList, UpdateCommand,
4};
5use teaql_data_service::{MutationRequest, QueryRequest};
6
7use crate::{DataServiceError, MetadataStore, RuntimeError};
8
9use super::RuntimeDataService;
10
11impl<'a, M, E> RuntimeDataService<'a, M, E>
12where
13    M: MetadataStore,
14    E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor,
15{
16    pub fn new(metadata: &'a M, executor: &'a E) -> Self {
17        Self { metadata, executor }
18    }
19
20    pub async fn fetch_all(
21        &self,
22        query: &SelectQuery,
23    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
24        let request = QueryRequest {
25            query: query.clone(),
26            trace_chain: query.trace_chain.clone(),
27            comment: query.comment.clone(),
28        };
29        let res = self
30            .executor
31            .query(request)
32            .await
33            .map_err(DataServiceError::Executor)?;
34        Ok(res.rows)
35    }
36
37    pub async fn fetch_smart_list(
38        &self,
39        query: &SelectQuery,
40    ) -> Result<SmartList<Record>, DataServiceError<E::Error>> {
41        let request = QueryRequest {
42            query: query.clone(),
43            trace_chain: query.trace_chain.clone(),
44            comment: query.comment.clone(),
45        };
46        let res = self
47            .executor
48            .query(request)
49            .await
50            .map_err(DataServiceError::Executor)?;
51        self.metadata.record_metadata_log(&res.metadata);
52        Ok(SmartList::from(res.rows))
53    }
54
55    pub async fn fetch_entities<T>(
56        &self,
57        query: &SelectQuery,
58    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
59    where
60        T: Entity,
61    {
62        self.fetch_all(query)
63            .await?
64            .into_iter()
65            .map(T::from_record)
66            .collect::<Result<Vec<_>, _>>()
67            .map(SmartList::from)
68            .map_err(DataServiceError::Entity)
69    }
70
71    pub async fn fetch_enhanced_entities<T>(
72        &self,
73        query: &SelectQuery,
74    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
75    where
76        T: Entity,
77    {
78        self.fetch_entities(query).await
79    }
80
81    pub async fn insert(&self, command: &InsertCommand) -> Result<u64, DataServiceError<E::Error>> {
82        let request = MutationRequest::Insert(command.clone());
83        let res = self
84            .executor
85            .mutate(request)
86            .await
87            .map_err(DataServiceError::Executor)?;
88        self.metadata.record_metadata_log(&res.metadata);
89        Ok(res.affected_rows)
90    }
91
92    pub async fn update(&self, command: &UpdateCommand) -> Result<u64, DataServiceError<E::Error>> {
93        let request = MutationRequest::Update(command.clone());
94        let res = self
95            .executor
96            .mutate(request)
97            .await
98            .map_err(DataServiceError::Executor)?;
99        self.metadata.record_metadata_log(&res.metadata);
100        let affected = res.affected_rows;
101
102        if command.expected_version.is_some() && affected == 0 {
103            println!(
104                "OptimisticLockConflict in base.rs update! entity={}, id={:?}",
105                command.entity, command.id
106            );
107            println!(
108                "Backtrace: {:#?}",
109                std::backtrace::Backtrace::force_capture()
110            );
111            return Err(DataServiceError::Runtime(
112                RuntimeError::OptimisticLockConflict {
113                    entity: command.entity.clone(),
114                    id: format!("{:?}", command.id),
115                },
116            ));
117        }
118
119        Ok(affected)
120    }
121
122    pub async fn delete(&self, command: &DeleteCommand) -> Result<u64, DataServiceError<E::Error>> {
123        let request = MutationRequest::Delete(command.clone());
124        let res = self
125            .executor
126            .mutate(request)
127            .await
128            .map_err(DataServiceError::Executor)?;
129        self.metadata.record_metadata_log(&res.metadata);
130        let affected = res.affected_rows;
131
132        if command.expected_version.is_some() && affected == 0 {
133            return Err(DataServiceError::Runtime(
134                RuntimeError::OptimisticLockConflict {
135                    entity: command.entity.clone(),
136                    id: format!("{:?}", command.id),
137                },
138            ));
139        }
140
141        Ok(affected)
142    }
143
144    pub async fn batch_insert(
145        &self,
146        command: &teaql_core::BatchInsertCommand,
147    ) -> Result<u64, DataServiceError<E::Error>> {
148        // Build individual InsertCommands for now, or use BatchMutation if appropriate
149        let mut affected = 0;
150        for (i, val) in command.batch_values.iter().enumerate() {
151            let mut insert_cmd = InsertCommand::new(command.entity.clone());
152            insert_cmd.values = val.clone();
153            if i < command.trace_chains.len() {
154                insert_cmd.trace_chain = command.trace_chains[i].clone();
155            }
156            let res = self
157                .executor
158                .mutate(MutationRequest::Insert(insert_cmd))
159                .await
160                .map_err(DataServiceError::Executor)?;
161            self.metadata.record_metadata_log(&res.metadata);
162            affected += res.affected_rows;
163        }
164        Ok(affected)
165    }
166
167    pub async fn batch_update(
168        &self,
169        command: &teaql_core::BatchUpdateCommand,
170    ) -> Result<u64, DataServiceError<E::Error>> {
171        let mut affected = 0;
172        for (i, val) in command.batch_values.iter().enumerate() {
173            let mut update_cmd =
174                UpdateCommand::new(command.entity.clone(), command.batch_ids[i].clone());
175
176            let mut filtered_values = Record::new();
177            for field in &command.update_fields {
178                if let Some(v) = val.get(field) {
179                    filtered_values.insert(field.clone(), v.clone());
180                }
181            }
182            update_cmd.values = filtered_values;
183            if let Some(Some(v)) = command.batch_expected_versions.get(i) {
184                update_cmd.expected_version = Some(*v);
185            }
186            if let Some(old) = command.batch_old_values.get(i) {
187                update_cmd.old_values = old.clone();
188            }
189            if i < command.trace_chains.len() {
190                update_cmd.trace_chain = command.trace_chains[i].clone();
191            }
192            let res = self
193                .executor
194                .mutate(MutationRequest::Update(update_cmd))
195                .await
196                .map_err(DataServiceError::Executor)?;
197            self.metadata.record_metadata_log(&res.metadata);
198            affected += res.affected_rows;
199        }
200
201        if command.batch_expected_versions.iter().any(|v| v.is_some()) {
202            if affected != command.batch_ids.len() as u64 {
203                println!(
204                    "OptimisticLockConflict in batch_update! entity={}, affected={}, expected={}",
205                    command.entity,
206                    affected,
207                    command.batch_ids.len()
208                );
209                return Err(DataServiceError::Runtime(
210                    RuntimeError::OptimisticLockConflict {
211                        entity: command.entity.clone(),
212                        id: "BATCH".to_owned(),
213                    },
214                ));
215            }
216        }
217
218        Ok(affected)
219    }
220
221    pub async fn recover(
222        &self,
223        command: &RecoverCommand,
224    ) -> Result<u64, DataServiceError<E::Error>> {
225        let request = MutationRequest::Recover(command.clone());
226        let res = self
227            .executor
228            .mutate(request)
229            .await
230            .map_err(DataServiceError::Executor)?;
231        self.metadata.record_metadata_log(&res.metadata);
232        let affected = res.affected_rows;
233
234        if affected == 0 {
235            return Err(DataServiceError::Runtime(
236                RuntimeError::OptimisticLockConflict {
237                    entity: command.entity.clone(),
238                    id: format!("{:?}", command.id),
239                },
240            ));
241        }
242
243        Ok(affected)
244    }
245
246    pub async fn insert_many(
247        &self,
248        commands: &[InsertCommand],
249    ) -> Result<u64, DataServiceError<E::Error>> {
250        let mut total = 0;
251        for command in commands {
252            total += self.insert(command).await?;
253        }
254        Ok(total)
255    }
256
257    pub async fn update_many(
258        &self,
259        commands: &[UpdateCommand],
260    ) -> Result<u64, DataServiceError<E::Error>> {
261        let mut total = 0;
262        for command in commands {
263            total += self.update(command).await?;
264        }
265        Ok(total)
266    }
267
268    pub async fn delete_many(
269        &self,
270        commands: &[DeleteCommand],
271    ) -> Result<u64, DataServiceError<E::Error>> {
272        let mut total = 0;
273        for command in commands {
274            total += self.delete(command).await?;
275        }
276        Ok(total)
277    }
278
279    pub async fn recover_many(
280        &self,
281        commands: &[RecoverCommand],
282    ) -> Result<u64, DataServiceError<E::Error>> {
283        let mut total = 0;
284        for command in commands {
285            total += self.recover(command).await?;
286        }
287        Ok(total)
288    }
289}