1use crate::cloud_outbox_records::CloudOutboxRecords;
2use crate::remote_object_records::load_remote_object_on;
3
4use super::*;
5
6pub(crate) fn carried_blob_locator(
13 remote: &coven_protocol::remote_object::RemoteObjectRecord,
14 context: &str,
15) -> Result<BlobLocator, DbError> {
16 let locator_bytes = remote.payloads().carried_locator_bytes().ok_or_else(|| {
17 DbError::Message(format!(
18 "{context}: stored blob {} carries no locator in its row",
19 remote.object_id()
20 ))
21 })?;
22 BlobLocator::parse(locator_bytes).map_err(|error| DbError::context(context.to_string(), error))
23}
24
25pub(crate) struct LiveBlobRow {
26 pub stamp: String,
27 pub blob_id: String,
28 pub plaintext_size: u64,
29 pub plaintext_hash: ObjectHash,
30 pub cloud_path: Option<String>,
31}
32
33pub(crate) fn live_blob_row(
34 conn: &Connection,
35 table: &str,
36 row_id: &str,
37 declaration: &coven_protocol::synced_schema::BlobDecl,
38) -> Result<Option<LiveBlobRow>, DbError> {
39 let cloud_path = declaration
40 .cloud_path_column
41 .as_deref()
42 .map(quote_ident)
43 .unwrap_or_else(|| "NULL".to_string());
44 let sql = format!(
45 "SELECT {}, {}, {}, {}, {} FROM {} WHERE {} = ?1",
46 quote_ident(&declaration.id_column),
47 quote_ident(&declaration.size_column),
48 quote_ident(&declaration.hash_column),
49 cloud_path,
50 quote_ident("_updated_at"),
51 quote_ident(table),
52 quote_ident("id"),
53 );
54 let raw = conn
55 .query_row(&sql, [row_id], |row| {
56 Ok((
57 row.get::<_, String>(0)?,
58 row.get::<_, i64>(1)?,
59 row.get::<_, String>(2)?,
60 row.get::<_, Option<String>>(3)?,
61 row.get::<_, String>(4)?,
62 ))
63 })
64 .optional()
65 .map_err(DbError::from)?;
66 let Some((blob_id, plaintext_size, plaintext_hash, cloud_path, stamp)) = raw else {
67 return Ok(None);
68 };
69 let plaintext_size = u64::try_from(plaintext_size).map_err(|_| {
70 DbError::Message(format!(
71 "winning blob row {:?}/{:?} has negative plaintext size {plaintext_size}",
72 table, row_id
73 ))
74 })?;
75 let plaintext_hash = plaintext_hash.parse().map_err(|error| {
76 DbError::context(
77 format!(
78 "winning blob row {:?}/{:?} has invalid plaintext hash",
79 table, row_id
80 ),
81 error,
82 )
83 })?;
84 Ok(Some(LiveBlobRow {
85 stamp,
86 blob_id,
87 plaintext_size,
88 plaintext_hash,
89 cloud_path,
90 }))
91}
92
93pub(crate) fn validate_live_blob_row(
94 binding: &RowBlobLocatorBinding,
95 declaration: &coven_protocol::synced_schema::BlobDecl,
96 row: &LiveBlobRow,
97 live_audience: &RemoteAudience,
98) -> Result<(), DbError> {
99 validate_live_blob_locator(
100 binding.table(),
101 binding.row_id(),
102 binding.column(),
103 binding.row_stamp(),
104 binding.blob(),
105 declaration,
106 row,
107 live_audience,
108 )
109}
110
111#[allow(clippy::too_many_arguments)]
112pub(crate) fn validate_live_blob_locator(
113 table: &str,
114 row_id: &str,
115 column: &str,
116 row_stamp: &str,
117 stored: &StoredBlobRef,
118 declaration: &coven_protocol::synced_schema::BlobDecl,
119 row: &LiveBlobRow,
120 live_audience: &RemoteAudience,
121) -> Result<(), DbError> {
122 let locator = stored.locator();
123 let invalid = locator.namespace() != declaration.namespace
124 || locator.blob_id() != row.blob_id
125 || locator.plaintext_size() != row.plaintext_size
126 || locator.plaintext_hash() != row.plaintext_hash
127 || &locator.audience() != live_audience
128 || locator
129 .scope()
130 .is_some_and(|scope| scope != &declaration.scope)
131 || locator
132 .cloud_path()
133 .is_some_and(|path| row.cloud_path.as_deref() != Some(path));
134 if invalid {
135 return Err(DbError::Message(format!(
136 "blob locator does not match winning row values for {:?}/{:?}/{:?} at {:?}",
137 table, row_id, column, row_stamp
138 )));
139 }
140 Ok(())
141}
142
143pub(crate) fn validate_stored_locator_on(
144 conn: &Connection,
145 expected: &StoredBlobRef,
146) -> Result<(), DbError> {
147 let locator_hash = expected.locator().locator_hash().to_string();
148 let expected_remote_object_id = remote_object_id(expected.object());
149 let stored_locator_hash: String = conn
150 .query_row(
151 "SELECT locator_hash FROM blob_locators WHERE remote_object_id = ?1",
152 [expected_remote_object_id.to_string()],
153 |row| row.get(0),
154 )
155 .map_err(DbError::from)?;
156 if stored_locator_hash != locator_hash {
157 return Err(DbError::Message(format!(
158 "stored blob object {expected_remote_object_id} is indexed under locator {stored_locator_hash}, expected {locator_hash}"
159 )));
160 }
161 let remote = load_remote_object_on(conn, expected_remote_object_id)?;
162 if !remote.is_activated_stored_blob() {
163 return Err(DbError::Message(format!(
164 "stored blob locator {locator_hash} does not reference activated ownership"
165 )));
166 }
167 let locator = carried_blob_locator(
168 &remote,
169 &format!("stored blob locator {locator_hash} is invalid"),
170 )?;
171 let actual = StoredBlobRef::new(locator, remote.object().clone()).map_err(|error| {
172 DbError::context(
173 format!("stored blob reference {locator_hash} is invalid"),
174 error,
175 )
176 })?;
177 if &actual != expected {
178 return Err(DbError::Message(format!(
179 "blob object {expected_remote_object_id} differs from its exact stored reference"
180 )));
181 }
182 Ok(())
183}
184
185pub(crate) fn validate_stored_row_binding_on(
186 conn: &Connection,
187 binding: &RowBlobLocatorBinding,
188 expected_authority: &coven_protocol::audience_package::PackageAudience,
189 expected_remote_object_id: ObjectHash,
190) -> Result<(), DbError> {
191 let (audience_authority, remote_object_id): (String, String) = conn
192 .query_row(
193 "SELECT audience_authority, remote_object_id FROM row_blob_locators
194 WHERE table_name = ?1 AND row_id = ?2 AND column_name = ?3 AND row_stamp = ?4",
195 rusqlite::params![
196 binding.table(),
197 binding.row_id(),
198 binding.column(),
199 binding.row_stamp(),
200 ],
201 |row| Ok((row.get(0)?, row.get(1)?)),
202 )
203 .map_err(DbError::from)?;
204 let actual_authority: coven_protocol::audience_package::PackageAudience =
205 serde_json::from_str(&audience_authority)
206 .map_err(|error| DbError::context("parse stored row blob audience authority", error))?;
207 if &actual_authority != expected_authority
208 || remote_object_id != expected_remote_object_id.to_string()
209 {
210 return Err(DbError::Message(format!(
211 "row blob binding {:?}/{:?}/{:?} at {:?} is already bound to different exact content",
212 binding.table(),
213 binding.row_id(),
214 binding.column(),
215 binding.row_stamp()
216 )));
217 }
218 Ok(())
219}
220
221pub(crate) fn load_prepared_audience_objects_on(
222 conn: &Connection,
223 store_dir: &coven_foundation::store_dir::StoreDir,
224 write_id: &WriteId,
225) -> Result<PreparedAudienceObjects, DbError> {
226 let mut package_statement = conn
227 .prepare(
228 "SELECT remote_object_id FROM store_write_packages
229 WHERE write_id = ?1 ORDER BY audience",
230 )
231 .map_err(DbError::from)?;
232 let package_ids = package_statement
233 .query_map([write_id.as_str()], |row| row.get::<_, String>(0))
234 .map_err(DbError::from)?
235 .collect::<Result<Vec<_>, _>>()
236 .map_err(DbError::from)?;
237 let mut blob_statement = conn
238 .prepare(
239 "SELECT remote_object_id, audience, locator_hash, spool_path FROM store_write_blobs
240 WHERE write_id = ?1 ORDER BY audience, remote_object_id",
241 )
242 .map_err(DbError::from)?;
243 let blob_rows = blob_statement
244 .query_map([write_id.as_str()], |row| {
245 Ok((
246 row.get::<_, String>(0)?,
247 row.get::<_, String>(1)?,
248 row.get::<_, String>(2)?,
249 row.get::<_, Option<String>>(3)?,
250 ))
251 })
252 .map_err(DbError::from)?
253 .collect::<Result<Vec<_>, _>>()
254 .map_err(DbError::from)?;
255 let packages = package_ids
256 .into_iter()
257 .map(|encoded| {
258 let object_id = encoded
259 .parse()
260 .map_err(|error| DbError::context("stored remote object id", error))?;
261 PreparedAudiencePackage::from_remote(
262 conn,
263 store_dir,
264 load_remote_object_on(conn, object_id)?,
265 )
266 })
267 .collect::<Result<Vec<_>, DbError>>()?;
268 let blobs = blob_rows
269 .into_iter()
270 .map(|(encoded, audience, locator_hash, spool_path)| {
271 let object_id = encoded
272 .parse()
273 .map_err(|error| DbError::context("stored remote object id", error))?;
274 PreparedAudienceBlob::from_remote(
275 parse_remote_audience_db(&audience)?,
276 &locator_hash,
277 load_remote_object_on(conn, object_id)?,
278 spool_path.map(PathBuf::from),
279 )
280 })
281 .collect::<Result<Vec<_>, DbError>>()?;
282 Ok(PreparedAudienceObjects { packages, blobs })
283}
284
285pub(crate) fn load_activated_registration_on(
286 conn: &Connection,
287 root: &coven_protocol::store_commit::StoreRootRef,
288 reference: &StoreDeviceRegistrationRef,
289) -> Result<StoreDeviceRegistration, DbError> {
290 let (bytes, encoded): (Vec<u8>, String) = conn
291 .query_row(
292 "SELECT registration_bytes, registration_object \
293 FROM store_device_registration_activations \
294 WHERE device_id = ?1 AND registration_hash = ?2",
295 (
296 reference.device_id.to_string(),
297 reference.registration_hash.to_string(),
298 ),
299 |row| Ok((row.get(0)?, row.get(1)?)),
300 )
301 .map_err(DbError::from)?;
302 let stored: StoreDeviceRegistrationRef = serde_json::from_str(&encoded)
303 .map_err(|error| DbError::context("activated Store registration ref", error))?;
304 if stored != *reference {
305 return Err(DbError::Message(
306 "activated Store registration differs from its exact reference".to_string(),
307 ));
308 }
309 let registration = StoreDeviceRegistration::parse_at(&bytes, root, reference.device_id)
310 .map_err(|error| DbError::context("activated Store registration", error))?;
311 reference
312 .verify_registration(®istration)
313 .map_err(DbError::from)?;
314 Ok(registration)
315}
316
317#[allow(clippy::too_many_arguments)]
318pub(crate) fn previous_row_blob_for_write_on(
319 conn: &Connection,
320 table: &str,
321 row_id: &str,
322 row_stamp: &str,
323 column: &str,
324 blob: &BlobRef,
325 plaintext_size: u64,
326 plaintext_hash: ObjectHash,
327) -> Result<Option<StoreWriteRemoteBlob>, DbError> {
328 if let Some(handoff) =
329 CloudOutboxRecords::new(conn).created_upload_handoff(table, row_id, column, row_stamp)?
330 {
331 let locator = handoff.stored.locator();
332 if !coven_protocol::blob::locator_describes_row(
333 locator,
334 blob,
335 plaintext_size,
336 plaintext_hash,
337 ) {
338 return Err(DbError::Message(format!(
339 "created upload {table}/{row_id}/{column} at {row_stamp} differs from its captured row"
340 )));
341 }
342 return Ok(Some(handoff));
343 }
344 let raw = conn
345 .query_row(
346 "SELECT row_blob_locators.audience_authority, blob_locators.remote_object_id
347 FROM row_blob_locators
348 JOIN blob_locators USING (remote_object_id)
349 WHERE table_name = ?1 AND row_id = ?2 AND column_name = ?3
350 ORDER BY row_stamp DESC LIMIT 1",
351 rusqlite::params![table, row_id, column],
352 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
353 )
354 .optional()
355 .map_err(DbError::from)?;
356 let Some((authority, object_id)) = raw else {
357 return Ok(None);
358 };
359 let authority: coven_protocol::audience_package::PackageAudience =
360 serde_json::from_str(&authority)
361 .map_err(|error| DbError::context("prior row blob authority", error))?;
362 let object_id = object_id
363 .parse()
364 .map_err(|error| DbError::context("prior row blob object id", error))?;
365 let remote = load_remote_object_on(conn, object_id)?;
366 if !remote.is_activated_stored_blob() {
367 return Err(DbError::Message(format!(
368 "prior row blob {table}/{row_id}/{column} is not activated"
369 )));
370 }
371 let locator = carried_blob_locator(&remote, "prior row blob locator")?;
372 if !coven_protocol::blob::locator_describes_row(&locator, blob, plaintext_size, plaintext_hash)
373 {
374 return Ok(None);
375 }
376 if locator.audience() != authority.remote_audience() {
377 return Err(DbError::Message(format!(
378 "prior row blob {table}/{row_id}/{column} authority differs from its locator"
379 )));
380 }
381 let stored = StoredBlobRef::new(locator, remote.object().clone())
382 .map_err(|error| DbError::context("prior row blob reference", error))?;
383 Ok(Some(StoreWriteRemoteBlob { authority, stored }))
384}
385
386pub fn remote_audience_to_db(audience: &RemoteAudience) -> String {
387 match audience {
388 RemoteAudience::Store => "store".to_string(),
389 RemoteAudience::Circle(circle_id) => circle_id.to_string(),
390 }
391}
392
393pub(crate) fn parse_remote_audience_db(value: &str) -> Result<RemoteAudience, DbError> {
394 if value == "store" {
395 return Ok(RemoteAudience::Store);
396 }
397 value
398 .parse()
399 .map(RemoteAudience::Circle)
400 .map_err(|error| DbError::context(format!("invalid stored blob audience {value:?}"), error))
401}