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