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 ¤t,
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!(¤t, 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}