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            .verified_prepared_store_commit_on(records, &prepared)?;
58        let mut ids = candidate_graph_exact_objects(commit.value())?
59            .iter()
60            .map(|object| (remote_object_id(object).to_string(), None))
61            .collect::<Vec<_>>();
62        let mut statement = self
63            .conn
64            .prepare(
65                "SELECT remote_object_id, spool_path
66                 FROM store_write_blobs WHERE write_id = ?1
67                 ORDER BY remote_object_id",
68            )
69            .map_err(DbError::from)?;
70        let blobs = statement
71            .query_map([write_id.as_str()], |row| {
72                Ok((row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?))
73            })
74            .map_err(DbError::from)?
75            .collect::<Result<Vec<_>, _>>()
76            .map_err(DbError::from)?;
77        ids.extend(blobs);
78        ids.sort_by(|left, right| left.0.cmp(&right.0));
79        ids.into_iter()
80            .map(|(encoded, spool_path)| {
81                let id = encoded
82                    .parse()
83                    .map_err(|error| DbError::context("prepared remote object id", error))?;
84                Ok(PreparedRemoteObject {
85                    closed: crate::reopen_remote_object_on(self.conn, self.store_dir, id)?,
86                    spool_path: spool_path.map(PathBuf::from),
87                })
88            })
89            .collect()
90    }
91
92    fn stage_membership_head_acceptance(
93        &self,
94        accepted: AcceptedStoreCommitEvidence,
95        closed: coven_protocol::remote_object::ClosedRemoteObject,
96    ) -> Result<RemoteObjectRecord, DbError> {
97        use coven_protocol::remote_object::{
98            RetainedAuthorityObjectDomain, RetainedAuthorityObjectState,
99        };
100        let RemoteObjectRecord::RetainedAuthority(proposed) = closed.record() else {
101            return Err(DbError::Message(
102                "membership acceptance is not retained authority".into(),
103            ));
104        };
105        let RetainedAuthorityObjectDomain::MembershipHeadAcceptance { publication, .. } =
106            &proposed.identity.domain
107        else {
108            return Err(DbError::Message(
109                "membership acceptance has another object domain".into(),
110            ));
111        };
112        let RetainedAuthorityObjectState::Prepared { ownership } = &proposed.state else {
113            return Err(DbError::Message(
114                "membership acceptance has no prepared owner".into(),
115            ));
116        };
117        if ownership.pending.len() != 1 || !ownership.pending.contains(accepted.commit_ref()) {
118            return Err(DbError::Message(
119                "membership acceptance belongs to another accepted commit".into(),
120            ));
121        }
122        if let Some(exact) = accepted.exact_publication() {
123            if exact.reference() != publication {
124                return Err(DbError::Message(
125                    "membership acceptance names another winning publication".into(),
126                ));
127            }
128        }
129        let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
130        let exists: bool = tx
131            .query_row(
132                "SELECT EXISTS(SELECT 1 FROM remote_objects WHERE object_id = ?1)",
133                [closed.object_id().to_string()],
134                |row| row.get(0),
135            )
136            .map_err(DbError::from)?;
137        let remote = if exists {
138            let existing = load_remote_object_on(&tx, closed.object_id())?;
139            let RemoteObjectRecord::RetainedAuthority(record) = &existing else {
140                return Err(DbError::Message(
141                    "membership acceptance changed ownership domain".into(),
142                ));
143            };
144            let owned = match &record.state {
145                RetainedAuthorityObjectState::Prepared { ownership } => {
146                    ownership.pending.contains(accepted.commit_ref())
147                }
148                RetainedAuthorityObjectState::UploadedVerified { ownership } => {
149                    ownership.pending.contains(accepted.commit_ref())
150                        || ownership.activated.contains(accepted.commit_ref())
151                }
152            };
153            if record.identity != proposed.identity
154                || record.payloads != proposed.payloads
155                || !owned
156            {
157                return Err(DbError::Message(
158                    "membership acceptance differs from its durable finalization owner".into(),
159                ));
160            }
161            existing
162        } else {
163            if accepted.exact_publication().is_none() {
164                return Err(DbError::Message(
165                    "covered membership acceptance has no durable finalization owner".into(),
166                ));
167            }
168            persist_exact_remote_object_on(
169                &tx,
170                self.store_dir,
171                &closed,
172                "accepted membership head result",
173            )?;
174            closed.record().clone()
175        };
176        tx.commit().map_err(DbError::from)?;
177        Ok(remote)
178    }
179
180    fn mark_remote_object_uploaded(
181        &self,
182        expected: RemoteObjectRecord,
183    ) -> Result<RemoteObjectRecord, DbError> {
184        mark_remote_object_uploaded_on(self.conn, expected)
185    }
186
187    fn uploaded_blob_spools(&self) -> Result<Vec<UploadedBlobSpool>, DbError> {
188        let conn = self.conn;
189        let mut statement = conn
190            .prepare(
191                "SELECT write_id, remote_object_id, spool_path
192                 FROM store_write_blobs
193                 WHERE spool_path IS NOT NULL
194                 ORDER BY write_id, audience, remote_object_id",
195            )
196            .map_err(DbError::from)?;
197        let rows = statement
198            .query_map([], |row| {
199                Ok((
200                    row.get::<_, String>(0)?,
201                    row.get::<_, String>(1)?,
202                    row.get::<_, String>(2)?,
203                ))
204            })
205            .map_err(DbError::from)?;
206        let rows = rows.collect::<Result<Vec<_>, _>>().map_err(DbError::from)?;
207        drop(statement);
208        let mut spools = Vec::new();
209        for (write_id, remote_object_id, path) in rows {
210            let remote_object_id = remote_object_id
211                .parse()
212                .map_err(|error| DbError::context("prepared blob remote object id", error))?;
213            let remote = load_remote_object_on(conn, remote_object_id)?;
214            if remote.records_verified_upload() {
215                spools.push(UploadedBlobSpool {
216                    write_id: WriteId::from_generated(write_id),
217                    remote_object_id,
218                    path: PathBuf::from(path),
219                });
220            }
221        }
222        Ok(spools)
223    }
224
225    fn clear_uploaded_blob_spool(&self, spool: UploadedBlobSpool) -> Result<(), DbError> {
226        let conn = self.conn;
227        let remote = load_remote_object_on(conn, spool.remote_object_id)?;
228        if !remote.records_verified_upload() {
229            return Err(DbError::Message(format!(
230                "prepared blob {} lost uploaded state before spool retirement",
231                spool.remote_object_id
232            )));
233        }
234        let path = spool
235            .path
236            .to_str()
237            .ok_or_else(|| DbError::Message("prepared blob spool path is not UTF-8".to_string()))?;
238        let cleared = conn
239            .execute(
240                "UPDATE store_write_blobs SET spool_path = NULL
241                 WHERE write_id = ?1 AND remote_object_id = ?2 AND spool_path = ?3",
242                rusqlite::params![
243                    spool.write_id.as_str(),
244                    spool.remote_object_id.to_string(),
245                    path,
246                ],
247            )
248            .map_err(DbError::from)?;
249        if cleared != 1 {
250            let current = conn
251                .query_row(
252                    "SELECT spool_path FROM store_write_blobs
253                     WHERE write_id = ?1 AND remote_object_id = ?2",
254                    rusqlite::params![spool.write_id.as_str(), spool.remote_object_id.to_string(),],
255                    |row| row.get::<_, Option<String>>(0),
256                )
257                .optional()
258                .map_err(DbError::from)?;
259            if current.flatten().is_some() {
260                return Err(DbError::Message(format!(
261                    "prepared blob {} spool changed during retirement",
262                    spool.remote_object_id
263                )));
264            }
265        }
266        Ok(())
267    }
268
269    fn mark_reusable_retained_authority_uploaded(
270        &self,
271        expected: RemoteObjectRecord,
272    ) -> Result<RemoteObjectRecord, DbError> {
273        mark_reusable_retained_authority_uploaded_on(self.conn, expected)
274    }
275
276    fn mark_candidate_commit_uploaded(&self, commit: StoreBatchCommitRef) -> Result<(), DbError> {
277        let conn = self.conn;
278        let object_id = remote_object_id(&commit.object);
279        let current = load_remote_object_on(conn, object_id)?;
280        if matches!(
281            &current,
282            RemoteObjectRecord::RetainedAuthority(record)
283                if matches!(
284                    &record.identity.domain,
285                    coven_protocol::remote_object::RetainedAuthorityObjectDomain::Commit {
286                        reference
287                    } if reference == &commit
288                ) && matches!(
289                    &record.state,
290                    coven_protocol::remote_object::RetainedAuthorityObjectState::UploadedVerified {
291                        ownership
292                    } if ownership.activated.contains(&commit)
293                )
294        ) {
295            return Ok(());
296        }
297        if !matches!(&current, RemoteObjectRecord::CandidateCommit(record) if record.identity == commit)
298        {
299            return Err(DbError::Message(format!(
300                "remote object {object_id} is not the exact candidate commit"
301            )));
302        }
303        mark_remote_object_uploaded_on(conn, current)?;
304        Ok(())
305    }
306
307    #[cfg(any(test, feature = "test-utils"))]
308    fn prepared_audience_objects(
309        &self,
310        write_id: &WriteId,
311    ) -> Result<PreparedAudienceObjects, DbError> {
312        load_prepared_audience_objects_on(self.conn, self.store_dir, write_id)
313    }
314
315    #[cfg(any(test, feature = "test-utils"))]
316    fn protocol_inert_object(
317        &self,
318        object: ExactObjectRef,
319    ) -> Result<Option<coven_protocol::remote_object::ProtocolInertObject>, DbError> {
320        let conn = self.conn;
321        let object_id = remote_object_id(&object);
322        let exists: bool = conn
323            .query_row(
324                "SELECT EXISTS(
325                    SELECT 1 FROM protocol_inert_objects WHERE object_id = ?1
326                 )",
327                [object_id.to_string()],
328                |row| row.get(0),
329            )
330            .map_err(DbError::from)?;
331        exists
332            .then(|| load_protocol_inert_object_on(conn, object_id))
333            .transpose()
334    }
335}
336
337impl StoreDatabase {
338    pub async fn prepared_remote_objects(
339        &self,
340        write_id: &WriteId,
341    ) -> Result<Vec<PreparedRemoteObject>, DbError> {
342        let write_id = write_id.clone();
343        self.call_store(move |session| session.prepared_remote_objects(&write_id))
344            .await
345    }
346
347    pub async fn stage_membership_head_acceptance(
348        &self,
349        accepted: AcceptedStoreCommitEvidence,
350        closed: coven_protocol::remote_object::ClosedRemoteObject,
351    ) -> Result<RemoteObjectRecord, DbError> {
352        self.call_store(move |session| session.stage_membership_head_acceptance(accepted, closed))
353            .await
354    }
355
356    pub async fn mark_remote_object_uploaded(
357        &self,
358        expected: RemoteObjectRecord,
359    ) -> Result<RemoteObjectRecord, DbError> {
360        self.call_store(move |session| session.mark_remote_object_uploaded(expected))
361            .await
362    }
363
364    pub async fn retire_uploaded_blob_spools(&self) -> Result<(), DbError> {
365        let spools = self
366            .call_store(|session| session.uploaded_blob_spools())
367            .await?;
368
369        for spool in spools {
370            spool.retire(self).await.map_err(DbError::File)?;
371            self.clear_uploaded_blob_spool(spool).await?;
372        }
373        Ok(())
374    }
375
376    async fn clear_uploaded_blob_spool(&self, spool: UploadedBlobSpool) -> Result<(), DbError> {
377        self.call_store(move |session| session.clear_uploaded_blob_spool(spool))
378            .await
379    }
380
381    pub async fn mark_reusable_retained_authority_uploaded(
382        &self,
383        expected: RemoteObjectRecord,
384    ) -> Result<RemoteObjectRecord, DbError> {
385        self.call_store(move |session| session.mark_reusable_retained_authority_uploaded(expected))
386            .await
387    }
388
389    pub async fn mark_candidate_commit_uploaded(
390        &self,
391        commit: StoreBatchCommitRef,
392    ) -> Result<(), DbError> {
393        self.call_store(move |session| session.mark_candidate_commit_uploaded(commit))
394            .await
395    }
396
397    #[cfg(any(test, feature = "test-utils"))]
398    pub async fn prepared_audience_objects(
399        &self,
400        write_id: &WriteId,
401    ) -> Result<PreparedAudienceObjects, DbError> {
402        let write_id = write_id.clone();
403        let loaded = self
404            .call_store(move |session| session.prepared_audience_objects(&write_id))
405            .await?;
406
407        let mut verified_blobs = Vec::with_capacity(loaded.blobs.len());
408        for prepared in loaded.blobs {
409            if let Some(spool_path) = prepared.spool_path() {
410                {
411                    let (size, digest) = coven_foundation::local_file::file_facts(spool_path)
412                        .await
413                        .map_err(DbError::File)?;
414                    prepared
415                        .blob()
416                        .object()
417                        .verify_stored_facts(
418                            spool_path,
419                            size,
420                            coven_protocol::store_commit::ObjectHash::from_digest(digest),
421                        )
422                        .map_err(|error| DbError::context("prepared blob spool", error))?;
423                }
424            }
425            verified_blobs.push(prepared);
426        }
427        Ok(PreparedAudienceObjects {
428            packages: loaded.packages,
429            blobs: verified_blobs,
430        })
431    }
432
433    #[cfg(any(test, feature = "test-utils"))]
434    pub async fn protocol_inert_object(
435        &self,
436        object: ExactObjectRef,
437    ) -> Result<Option<coven_protocol::remote_object::ProtocolInertObject>, DbError> {
438        self.call_store(move |session| session.protocol_inert_object(object))
439            .await
440    }
441}
442
443pub(crate) fn persist_prepared_audience_objects_on(
444    conn: &rusqlite::Transaction<'_>,
445    store_dir: &coven_foundation::store_dir::StoreDir,
446    write_id: &WriteId,
447    packages: &[PreparedAudiencePackage],
448    blobs: &[PreparedAudienceBlob],
449) -> Result<(), DbError> {
450    let package_audiences = packages
451        .iter()
452        .map(|prepared| {
453            if prepared.package().write_id() != write_id {
454                return Err(DbError::Message(format!(
455                    "prepared audience package write {} differs from journal write {write_id}",
456                    prepared.package().write_id()
457                )));
458            }
459            Ok(prepared.package().audience().remote_audience())
460        })
461        .collect::<Result<std::collections::BTreeSet<_>, DbError>>()?;
462    if package_audiences.len() != packages.len() {
463        return Err(DbError::Message(format!(
464            "write {write_id} has duplicate prepared package audiences"
465        )));
466    }
467    for prepared in packages {
468        let audience = prepared.package().audience().remote_audience();
469        validate_remote_object_on(
470            conn,
471            prepared.remote_object_id(),
472            prepared.object(),
473            prepared.semantic_bytes(),
474        )?;
475        conn.execute(
476            "INSERT INTO store_write_packages
477             (write_id, audience, remote_object_id)
478             VALUES (?1, ?2, ?3)
479             ON CONFLICT(write_id, audience) DO NOTHING",
480            rusqlite::params![
481                write_id.as_str(),
482                remote_audience_to_db(&audience),
483                prepared.remote_object_id().to_string(),
484            ],
485        )
486        .map_err(DbError::from)?;
487        validate_prepared_package_on(conn, store_dir, write_id, prepared)?;
488    }
489    for prepared in blobs {
490        if !package_audiences.contains(prepared.audience()) {
491            return Err(DbError::Message(format!(
492                "write {write_id} has a prepared blob for {:?} without that audience's package",
493                prepared.audience()
494            )));
495        }
496        let locator = prepared.blob().locator();
497        validate_remote_object_on(
498            conn,
499            prepared.remote_object_id(),
500            prepared.blob().object(),
501            &locator.to_bytes(),
502        )?;
503        let locator_hash = locator.locator_hash();
504        let spool_path = prepared
505            .spool_path()
506            .map(|path| {
507                path.to_str().map(str::to_string).ok_or_else(|| {
508                    DbError::Message("prepared blob spool path is not UTF-8".to_string())
509                })
510            })
511            .transpose()?;
512        conn.execute(
513            "INSERT INTO store_write_blobs
514             (write_id, audience, locator_hash, remote_object_id, spool_path)
515             VALUES (?1, ?2, ?3, ?4, ?5)
516             ON CONFLICT(write_id, audience, remote_object_id) DO NOTHING",
517            rusqlite::params![
518                write_id.as_str(),
519                remote_audience_to_db(prepared.audience()),
520                locator_hash.to_string(),
521                prepared.remote_object_id().to_string(),
522                spool_path,
523            ],
524        )
525        .map_err(DbError::from)?;
526        validate_prepared_blob_on(conn, write_id, prepared)?;
527    }
528    Ok(())
529}