teaql_runtime/data_service/
base.rs1use 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 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}