eidetica/instance/backend/
remote.rs1use std::{
4 collections::BTreeMap,
5 sync::{Arc, Mutex},
6};
7
8use async_trait::async_trait;
9
10use super::{Backend, MergeSlice};
11use crate::{
12 Result,
13 auth::SigKey,
14 backend::{
15 InstanceMetadata, RecordMutations, RecordPage, RecordRange, RecordView, StagingToken,
16 StoreStateRequest, VerificationStatus,
17 },
18 entry::{Entry, ID},
19 instance::WriteSource,
20 service::{client::RemoteConnection, protocol::ReadScope},
21 snapshot::Snapshot,
22};
23
24#[derive(Debug, Clone)]
43pub struct RemoteBackend {
44 conn: RemoteConnection,
45 identity: Option<SigKey>,
46 views: Arc<Mutex<BTreeMap<String, ID>>>,
47}
48
49impl RemoteBackend {
50 pub fn new(conn: RemoteConnection, identity: Option<SigKey>) -> Self {
51 Self {
52 conn,
53 identity,
54 views: Arc::new(Mutex::new(BTreeMap::new())),
55 }
56 }
57
58 fn identity(&self) -> SigKey {
61 self.identity
62 .clone()
63 .or_else(|| self.conn.session_identity())
64 .unwrap_or_default()
65 }
66}
67
68#[async_trait]
69impl Backend for RemoteBackend {
70 async fn resolve_store_state(&self, request: &StoreStateRequest) -> Result<Option<RecordView>> {
71 let view = self
72 .conn
73 .resolve_store_state(self.identity(), request.clone())
74 .await?;
75 if let Some(token) = &view {
76 self.views
77 .lock()
78 .unwrap()
79 .insert(token.clone(), request.database.clone());
80 }
81 Ok(view.map(|namespace_id| RecordView { namespace_id }))
82 }
83
84 async fn begin_store_state_staging(&self, request: StoreStateRequest) -> Result<StagingToken> {
85 let namespace_id = self
86 .conn
87 .begin_store_state_staging(self.identity(), request.clone())
88 .await?;
89 Ok(StagingToken {
90 namespace_id,
91 target: request,
92 })
93 }
94
95 async fn stage_store_state_records(
96 &self,
97 token: &StagingToken,
98 records: RecordMutations,
99 ) -> Result<()> {
100 let mut chunk_id = 0;
101 let mut chunk = RecordMutations::new();
102 let mut encoded = 0usize;
103 for (key, value) in records {
104 let size = serde_json::to_vec(&(key.clone(), value.clone()))?.len();
105 if size > crate::service::protocol::MAX_RECORD_CHUNK_BYTES as usize {
106 return Err(crate::backend::BackendError::RecordTooLarge {
107 encoded_bytes: size,
108 }
109 .into());
110 }
111 if !chunk.is_empty()
112 && encoded + size > crate::service::protocol::MAX_RECORD_CHUNK_BYTES as usize
113 {
114 self.conn
115 .stage_store_state_records(
116 token.target.database.clone(),
117 self.identity(),
118 token.namespace_id.clone(),
119 chunk_id,
120 std::mem::take(&mut chunk),
121 )
122 .await?;
123 chunk_id += 1;
124 encoded = 0;
125 }
126 encoded += size;
127 chunk.insert(key, value);
128 }
129 if !chunk.is_empty() {
130 self.conn
131 .stage_store_state_records(
132 token.target.database.clone(),
133 self.identity(),
134 token.namespace_id.clone(),
135 chunk_id,
136 chunk,
137 )
138 .await?;
139 }
140 Ok(())
141 }
142
143 async fn publish_store_state(&self, token: StagingToken) -> Result<RecordView> {
144 let database = token.target.database.clone();
145 let namespace_id = self
146 .conn
147 .publish_store_state(token.target.database, self.identity(), token.namespace_id)
148 .await?;
149 self.views
150 .lock()
151 .unwrap()
152 .insert(namespace_id.clone(), database);
153 Ok(RecordView { namespace_id })
154 }
155
156 async fn abort_store_state(&self, token: StagingToken) -> Result<()> {
157 self.conn
158 .abort_store_state(token.target.database, self.identity(), token.namespace_id)
159 .await
160 }
161
162 async fn store_state_record_get(
163 &self,
164 view: &RecordView,
165 key: &[u8],
166 ) -> Result<Option<Vec<u8>>> {
167 let database = self
168 .views
169 .lock()
170 .unwrap()
171 .get(&view.namespace_id)
172 .cloned()
173 .ok_or(crate::backend::BackendError::InvalidStoreStateView)?;
174 self.conn
175 .store_state_record_get(
176 database,
177 self.identity(),
178 view.namespace_id.clone(),
179 key.to_vec(),
180 )
181 .await
182 }
183
184 async fn store_state_record_scan(
185 &self,
186 view: &RecordView,
187 range: &RecordRange,
188 after: Option<&[u8]>,
189 limit: usize,
190 ) -> Result<RecordPage> {
191 let database = self
192 .views
193 .lock()
194 .unwrap()
195 .get(&view.namespace_id)
196 .cloned()
197 .ok_or(crate::backend::BackendError::InvalidStoreStateView)?;
198 self.conn
199 .store_state_record_scan(
200 database,
201 self.identity(),
202 view.namespace_id.clone(),
203 range.clone(),
204 after.map(ToOwned::to_owned),
205 u32::try_from(limit).unwrap_or(u32::MAX),
206 crate::service::protocol::MAX_RECORD_PAGE_BYTES,
207 )
208 .await
209 }
210
211 async fn clear_derived_store_state(&self) -> Result<()> {
212 Ok(())
217 }
218
219 async fn get(&self, id: &ID) -> Result<Entry> {
220 self.conn
224 .db_get_entry(ID::default(), self.identity(), id.clone())
225 .await
226 }
227
228 async fn snapshot(&self, tree: &ID) -> Result<Snapshot> {
229 match self
230 .conn
231 .get_verified_tips(tree.clone(), self.identity())
232 .await
233 {
234 Ok(snapshot) => Ok(snapshot),
235 Err(e) if e.is_not_found() => Ok(Snapshot::EMPTY),
236 Err(e) => Err(e),
237 }
238 }
239
240 async fn store_snapshot(&self, tree: &ID, store: &str) -> Result<Snapshot> {
241 let tree_tips = match self
242 .conn
243 .get_verified_tips(tree.clone(), self.identity())
244 .await
245 {
246 Ok(tips) => tips,
247 Err(e) if e.is_not_found() => return Ok(Snapshot::EMPTY),
248 Err(e) => return Err(e),
249 };
250 if tree_tips.is_empty() {
251 return Ok(Snapshot::EMPTY);
252 }
253 match self
254 .conn
255 .store_snapshot_at(
256 tree.clone(),
257 self.identity(),
258 store.to_string(),
259 tree_tips.into_tips(),
260 )
261 .await
262 {
263 Ok(snapshot) => Ok(snapshot),
264 Err(e) if e.is_not_found() => Ok(Snapshot::EMPTY),
265 Err(e) => Err(e),
266 }
267 }
268
269 async fn store_snapshot_at(
270 &self,
271 tree: &ID,
272 store: &str,
273 main_snapshot: &Snapshot,
274 ) -> Result<Snapshot> {
275 match self
276 .conn
277 .store_snapshot_at(
278 tree.clone(),
279 self.identity(),
280 store.to_string(),
281 main_snapshot.tips().to_vec(),
282 )
283 .await
284 {
285 Ok(snapshot) => Ok(snapshot),
286 Err(e) if e.is_not_found() => Ok(Snapshot::EMPTY),
287 Err(e) => Err(e),
288 }
289 }
290
291 async fn store_at(&self, tree: &ID, store: &str, snapshot: &Snapshot) -> Result<Vec<Entry>> {
292 self.conn
293 .get_store_entries(
294 tree.clone(),
295 self.identity(),
296 store.to_string(),
297 snapshot.tips().to_vec(),
298 ReadScope::Verified,
299 )
300 .await
301 }
302
303 async fn compute_merge_state(
304 &self,
305 tree: &ID,
306 store: &str,
307 entry_ids: &[ID],
308 ) -> Result<MergeSlice> {
309 let state = self
312 .conn
313 .compute_merge_state(
314 tree.clone(),
315 self.identity(),
316 store.to_string(),
317 entry_ids.to_vec(),
318 )
319 .await?;
320 Ok(MergeSlice {
321 merge_base: state.merge_base,
322 path: state.path,
323 })
324 }
325
326 async fn put(&self, entry: Entry) -> Result<()> {
327 let tree_root = entry.root().unwrap_or_else(|| entry.id());
328 self.conn
329 .submit_signed_entry(tree_root, self.identity(), entry)
330 .await
331 }
332
333 async fn write_entry(
334 &self,
335 _verification: VerificationStatus,
336 entry: Entry,
337 _source: WriteSource,
338 ) -> Result<()> {
339 let tree_root = entry.root().unwrap_or_else(|| entry.id());
342 self.conn
343 .submit_signed_entry(tree_root, self.identity(), entry)
344 .await
345 }
346
347 async fn get_instance_metadata(&self) -> Result<Option<InstanceMetadata>> {
348 self.conn.get_instance_metadata().await
349 }
350
351 async fn set_instance_metadata(&self, metadata: &InstanceMetadata) -> Result<()> {
352 self.conn.set_instance_metadata(metadata).await
353 }
354
355 fn remote_connection(&self) -> Option<RemoteConnection> {
356 Some(self.conn.clone())
357 }
358}