Skip to main content

coven_database/store/store_session/
prepared_remote_objects.rs

1use crate::*;
2#[cfg(any(test, feature = "test-utils"))]
3use coven_protocol::objects::ExactObjectRef;
4use coven_protocol::remote_object::{remote_object_id, RemoteObjectRecord};
5use coven_protocol::store_commit::{ObjectHash, StoreBatchCommitRef};
6use coven_protocol::write::WriteId;
7use rusqlite::OptionalExtension;
8use std::path::PathBuf;
9
10use super::publication_state::PreparedStoreWriteState;
11use super::*;
12
13struct UploadedBlobSpool {
14    write_id: WriteId,
15    remote_object_id: ObjectHash,
16    path: PathBuf,
17}
18
19impl UploadedBlobSpool {
20    async fn retire(
21        &self,
22        database: &StoreDatabase,
23    ) -> Result<(), coven_foundation::atomic_file::FileError> {
24        match tokio::fs::remove_file(&self.path).await {
25            Ok(()) => {}
26            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
27            Err(source) => {
28                return Err(coven_foundation::atomic_file::FileError::Path {
29                    operation: "remove uploaded prepared blob spool",
30                    path: self.path.clone(),
31                    source,
32                })
33            }
34        }
35        database.sync_store_parent_dir(&self.path).await
36    }
37}
38
39impl StoreSession<'_> {
40    fn prepared_remote_objects(
41        &mut self,
42        write_id: &WriteId,
43    ) -> Result<Vec<PreparedRemoteObject>, DbError> {
44        let records = crate::store::store_session::StoreRecords::new(self.conn, self.store_dir);
45        let raw_prepared: String = self
46            .conn
47            .query_row(
48                "SELECT prepared FROM store_writes WHERE write_id = ?1",
49                [write_id.as_str()],
50                |row| row.get(0),
51            )
52            .map_err(DbError::from)?;
53        let prepared: PreparedStoreWriteState = serde_json::from_str(&raw_prepared)
54            .map_err(|error| DbError::context("prepared remote graph", error))?;
55        let commit = self
56            .verified_store_authority
57            .prepared_merge_candidate_on(records, &prepared)?
58            .commit;
59        let mut ids = candidate_graph_exact_objects(&commit)?
60            .iter()
61            .map(|object| (remote_object_id(object).to_string(), None))
62            .collect::<Vec<_>>();
63        let mut statement = self
64            .conn
65            .prepare(
66                "SELECT remote_object_id, spool_path
67                 FROM store_write_blobs WHERE write_id = ?1
68                 ORDER BY remote_object_id",
69            )
70            .map_err(DbError::from)?;
71        let blobs = statement
72            .query_map([write_id.as_str()], |row| {
73                Ok((row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?))
74            })
75            .map_err(DbError::from)?
76            .collect::<Result<Vec<_>, _>>()
77            .map_err(DbError::from)?;
78        ids.extend(blobs);
79        ids.sort_by(|left, right| left.0.cmp(&right.0));
80        ids.into_iter()
81            .map(|(encoded, spool_path)| {
82                let id = encoded
83                    .parse()
84                    .map_err(|error| DbError::context("prepared remote object id", error))?;
85                Ok(PreparedRemoteObject {
86                    closed: crate::reopen_remote_object_on(self.conn, self.store_dir, id)?,
87                    spool_path: spool_path.map(PathBuf::from),
88                })
89            })
90            .collect()
91    }
92
93    fn mark_remote_object_uploaded(
94        &self,
95        expected: RemoteObjectRecord,
96    ) -> Result<RemoteObjectRecord, DbError> {
97        mark_remote_object_uploaded_on(self.conn, expected)
98    }
99
100    fn uploaded_blob_spools(&self) -> Result<Vec<UploadedBlobSpool>, DbError> {
101        let conn = self.conn;
102        let mut statement = conn
103            .prepare(
104                "SELECT write_id, remote_object_id, spool_path
105                 FROM store_write_blobs
106                 WHERE spool_path IS NOT NULL
107                 ORDER BY write_id, audience, remote_object_id",
108            )
109            .map_err(DbError::from)?;
110        let rows = statement
111            .query_map([], |row| {
112                Ok((
113                    row.get::<_, String>(0)?,
114                    row.get::<_, String>(1)?,
115                    row.get::<_, String>(2)?,
116                ))
117            })
118            .map_err(DbError::from)?;
119        let rows = rows.collect::<Result<Vec<_>, _>>().map_err(DbError::from)?;
120        drop(statement);
121        let mut spools = Vec::new();
122        for (write_id, remote_object_id, path) in rows {
123            let remote_object_id = remote_object_id
124                .parse()
125                .map_err(|error| DbError::context("prepared blob remote object id", error))?;
126            let remote = load_remote_object_on(conn, remote_object_id)?;
127            if remote.records_verified_upload() {
128                spools.push(UploadedBlobSpool {
129                    write_id: WriteId::from_generated(write_id),
130                    remote_object_id,
131                    path: PathBuf::from(path),
132                });
133            }
134        }
135        Ok(spools)
136    }
137
138    fn clear_uploaded_blob_spool(&self, spool: UploadedBlobSpool) -> Result<(), DbError> {
139        let conn = self.conn;
140        let remote = load_remote_object_on(conn, spool.remote_object_id)?;
141        if !remote.records_verified_upload() {
142            return Err(DbError::Message(format!(
143                "prepared blob {} lost uploaded state before spool retirement",
144                spool.remote_object_id
145            )));
146        }
147        let path = spool
148            .path
149            .to_str()
150            .ok_or_else(|| DbError::Message("prepared blob spool path is not UTF-8".to_string()))?;
151        let cleared = conn
152            .execute(
153                "UPDATE store_write_blobs SET spool_path = NULL
154                 WHERE write_id = ?1 AND remote_object_id = ?2 AND spool_path = ?3",
155                rusqlite::params![
156                    spool.write_id.as_str(),
157                    spool.remote_object_id.to_string(),
158                    path,
159                ],
160            )
161            .map_err(DbError::from)?;
162        if cleared != 1 {
163            let current = conn
164                .query_row(
165                    "SELECT spool_path FROM store_write_blobs
166                     WHERE write_id = ?1 AND remote_object_id = ?2",
167                    rusqlite::params![spool.write_id.as_str(), spool.remote_object_id.to_string(),],
168                    |row| row.get::<_, Option<String>>(0),
169                )
170                .optional()
171                .map_err(DbError::from)?;
172            if current.flatten().is_some() {
173                return Err(DbError::Message(format!(
174                    "prepared blob {} spool changed during retirement",
175                    spool.remote_object_id
176                )));
177            }
178        }
179        Ok(())
180    }
181
182    fn mark_reusable_retained_authority_uploaded(
183        &self,
184        expected: RemoteObjectRecord,
185    ) -> Result<RemoteObjectRecord, DbError> {
186        mark_reusable_retained_authority_uploaded_on(self.conn, expected)
187    }
188
189    fn mark_candidate_commit_uploaded(&self, commit: StoreBatchCommitRef) -> Result<(), DbError> {
190        let conn = self.conn;
191        let object_id = remote_object_id(&commit.object);
192        let current = load_remote_object_on(conn, object_id)?;
193        if matches!(
194            &current,
195            RemoteObjectRecord::RetainedAuthority(record)
196                if matches!(
197                    &record.identity.domain,
198                    coven_protocol::remote_object::RetainedAuthorityObjectDomain::Commit {
199                        reference
200                    } if reference == &commit
201                ) && matches!(
202                    &record.state,
203                    coven_protocol::remote_object::RetainedAuthorityObjectState::UploadedVerified {
204                        ownership
205                    } if ownership.activated.contains(&commit)
206                )
207        ) {
208            return Ok(());
209        }
210        if !matches!(&current, RemoteObjectRecord::CandidateCommit(record) if record.identity == commit)
211        {
212            return Err(DbError::Message(format!(
213                "remote object {object_id} is not the exact candidate commit"
214            )));
215        }
216        mark_remote_object_uploaded_on(conn, current)?;
217        Ok(())
218    }
219
220    fn mark_store_head_uploaded(
221        &self,
222        head: coven_protocol::store_commit::StoreDeviceHeadRef,
223    ) -> Result<(), DbError> {
224        let conn = self.conn;
225        let object_id = remote_object_id(&head.object);
226        let current = load_remote_object_on(conn, object_id)?;
227        if !matches!(
228            &current,
229            RemoteObjectRecord::RetainedAuthority(record)
230                if matches!(
231                    &record.identity.domain,
232                    coven_protocol::remote_object::RetainedAuthorityObjectDomain::DeviceHead {
233                        reference,
234                        ..
235                    } if reference == &head
236                )
237        ) {
238            return Err(DbError::Message(format!(
239                "remote object {object_id} is not the exact prepared Store head"
240            )));
241        }
242        mark_remote_object_uploaded_on(conn, current)?;
243        Ok(())
244    }
245
246    #[cfg(any(test, feature = "test-utils"))]
247    fn prepared_audience_objects(
248        &self,
249        write_id: &WriteId,
250    ) -> Result<PreparedAudienceObjects, DbError> {
251        load_prepared_audience_objects_on(self.conn, self.store_dir, write_id)
252    }
253
254    #[cfg(any(test, feature = "test-utils"))]
255    fn protocol_inert_object(
256        &self,
257        object: ExactObjectRef,
258    ) -> Result<Option<coven_protocol::remote_object::ProtocolInertObject>, DbError> {
259        let conn = self.conn;
260        let object_id = remote_object_id(&object);
261        let exists: bool = conn
262            .query_row(
263                "SELECT EXISTS(
264                    SELECT 1 FROM protocol_inert_objects WHERE object_id = ?1
265                 )",
266                [object_id.to_string()],
267                |row| row.get(0),
268            )
269            .map_err(DbError::from)?;
270        exists
271            .then(|| load_protocol_inert_object_on(conn, object_id))
272            .transpose()
273    }
274}
275
276impl StoreDatabase {
277    pub async fn prepared_remote_objects(
278        &self,
279        write_id: &WriteId,
280    ) -> Result<Vec<PreparedRemoteObject>, DbError> {
281        let write_id = write_id.clone();
282        self.call_store(move |session| session.prepared_remote_objects(&write_id))
283            .await
284    }
285
286    pub async fn mark_remote_object_uploaded(
287        &self,
288        expected: RemoteObjectRecord,
289    ) -> Result<RemoteObjectRecord, DbError> {
290        self.call_store(move |session| session.mark_remote_object_uploaded(expected))
291            .await
292    }
293
294    pub async fn retire_uploaded_blob_spools(&self) -> Result<(), DbError> {
295        let spools = self
296            .call_store(|session| session.uploaded_blob_spools())
297            .await?;
298
299        for spool in spools {
300            spool.retire(self).await.map_err(DbError::File)?;
301            self.clear_uploaded_blob_spool(spool).await?;
302        }
303        Ok(())
304    }
305
306    async fn clear_uploaded_blob_spool(&self, spool: UploadedBlobSpool) -> Result<(), DbError> {
307        self.call_store(move |session| session.clear_uploaded_blob_spool(spool))
308            .await
309    }
310
311    pub async fn mark_reusable_retained_authority_uploaded(
312        &self,
313        expected: RemoteObjectRecord,
314    ) -> Result<RemoteObjectRecord, DbError> {
315        self.call_store(move |session| session.mark_reusable_retained_authority_uploaded(expected))
316            .await
317    }
318
319    pub async fn mark_candidate_commit_uploaded(
320        &self,
321        commit: StoreBatchCommitRef,
322    ) -> Result<(), DbError> {
323        self.call_store(move |session| session.mark_candidate_commit_uploaded(commit))
324            .await
325    }
326
327    pub async fn mark_store_head_uploaded(
328        &self,
329        head: coven_protocol::store_commit::StoreDeviceHeadRef,
330    ) -> Result<(), DbError> {
331        self.call_store(move |session| session.mark_store_head_uploaded(head))
332            .await
333    }
334
335    #[cfg(any(test, feature = "test-utils"))]
336    pub async fn prepared_audience_objects(
337        &self,
338        write_id: &WriteId,
339    ) -> Result<PreparedAudienceObjects, DbError> {
340        let write_id = write_id.clone();
341        let loaded = self
342            .call_store(move |session| session.prepared_audience_objects(&write_id))
343            .await?;
344
345        let mut verified_blobs = Vec::with_capacity(loaded.blobs.len());
346        for prepared in loaded.blobs {
347            if let Some(spool_path) = prepared.spool_path() {
348                {
349                    let (size, digest) = coven_foundation::local_file::file_facts(spool_path)
350                        .await
351                        .map_err(DbError::File)?;
352                    prepared
353                        .blob()
354                        .object()
355                        .verify_stored_facts(
356                            spool_path,
357                            size,
358                            coven_protocol::store_commit::ObjectHash::from_digest(digest),
359                        )
360                        .map_err(|error| DbError::context("prepared blob spool", error))?;
361                }
362            }
363            verified_blobs.push(prepared);
364        }
365        Ok(PreparedAudienceObjects {
366            packages: loaded.packages,
367            blobs: verified_blobs,
368        })
369    }
370
371    #[cfg(any(test, feature = "test-utils"))]
372    pub async fn protocol_inert_object(
373        &self,
374        object: ExactObjectRef,
375    ) -> Result<Option<coven_protocol::remote_object::ProtocolInertObject>, DbError> {
376        self.call_store(move |session| session.protocol_inert_object(object))
377            .await
378    }
379}
380
381pub(crate) fn persist_prepared_audience_objects_on(
382    conn: &rusqlite::Transaction<'_>,
383    store_dir: &coven_foundation::store_dir::StoreDir,
384    write_id: &WriteId,
385    packages: &[PreparedAudiencePackage],
386    blobs: &[PreparedAudienceBlob],
387) -> Result<(), DbError> {
388    let package_audiences = packages
389        .iter()
390        .map(|prepared| {
391            if prepared.package().write_id() != write_id {
392                return Err(DbError::Message(format!(
393                    "prepared audience package write {} differs from journal write {write_id}",
394                    prepared.package().write_id()
395                )));
396            }
397            Ok(prepared.package().audience().remote_audience())
398        })
399        .collect::<Result<std::collections::BTreeSet<_>, DbError>>()?;
400    if package_audiences.len() != packages.len() {
401        return Err(DbError::Message(format!(
402            "write {write_id} has duplicate prepared package audiences"
403        )));
404    }
405    for prepared in packages {
406        let audience = prepared.package().audience().remote_audience();
407        validate_remote_object_on(
408            conn,
409            prepared.remote_object_id(),
410            prepared.object(),
411            prepared.semantic_bytes(),
412        )?;
413        conn.execute(
414            "INSERT INTO store_write_packages
415             (write_id, audience, remote_object_id)
416             VALUES (?1, ?2, ?3)
417             ON CONFLICT(write_id, audience) DO NOTHING",
418            rusqlite::params![
419                write_id.as_str(),
420                remote_audience_to_db(&audience),
421                prepared.remote_object_id().to_string(),
422            ],
423        )
424        .map_err(DbError::from)?;
425        validate_prepared_package_on(conn, store_dir, write_id, prepared)?;
426    }
427    for prepared in blobs {
428        if !package_audiences.contains(prepared.audience()) {
429            return Err(DbError::Message(format!(
430                "write {write_id} has a prepared blob for {:?} without that audience's package",
431                prepared.audience()
432            )));
433        }
434        let locator = prepared.blob().locator();
435        validate_remote_object_on(
436            conn,
437            prepared.remote_object_id(),
438            prepared.blob().object(),
439            &locator.to_bytes(),
440        )?;
441        let locator_hash = locator.locator_hash();
442        let spool_path = prepared
443            .spool_path()
444            .map(|path| {
445                path.to_str().map(str::to_string).ok_or_else(|| {
446                    DbError::Message("prepared blob spool path is not UTF-8".to_string())
447                })
448            })
449            .transpose()?;
450        conn.execute(
451            "INSERT INTO store_write_blobs
452             (write_id, audience, locator_hash, remote_object_id, spool_path)
453             VALUES (?1, ?2, ?3, ?4, ?5)
454             ON CONFLICT(write_id, audience, remote_object_id) DO NOTHING",
455            rusqlite::params![
456                write_id.as_str(),
457                remote_audience_to_db(prepared.audience()),
458                locator_hash.to_string(),
459                prepared.remote_object_id().to_string(),
460                spool_path,
461            ],
462        )
463        .map_err(DbError::from)?;
464        validate_prepared_blob_on(conn, write_id, prepared)?;
465    }
466    Ok(())
467}