1use super::*;
2
3#[derive(Debug, Clone, PartialEq, Eq)]
5pub struct BlobActivation {
6 pub coord: StoreCommitCoord,
7}
8
9#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct PreparedAudiencePackage {
11 remote_object_id: ObjectHash,
12 package: AudiencePackage,
13 semantic_bytes: Vec<u8>,
14 stored_bytes: Vec<u8>,
15 object: ExactObjectRef,
16}
17
18impl PreparedAudiencePackage {
19 pub(crate) fn from_remote(
22 conn: &rusqlite::Connection,
23 store_dir: &coven_foundation::store_dir::StoreDir,
24 remote: RemoteObjectRecord,
25 ) -> Result<Self, DbError> {
26 remote
27 .validate()
28 .map_err(|error| DbError::context("prepared remote package", error))?;
29 let is_package = match &remote {
30 RemoteObjectRecord::CandidateCommit(_) | RemoteObjectRecord::RetainedAuthority(_) => {
31 false
32 }
33 RemoteObjectRecord::CandidateExclusive(record) => matches!(
34 record.identity.domain,
35 CandidateExclusiveObjectDomain::StorePackage { .. }
36 | CandidateExclusiveObjectDomain::CirclePackage { .. }
37 ),
38 RemoteObjectRecord::SharedLiveSet(record) => matches!(
39 record.identity.domain,
40 SharedLiveSetObjectDomain::StorePackage { .. }
41 | SharedLiveSetObjectDomain::CirclePackage { .. }
42 ),
43 };
44 if !is_package {
45 return Err(DbError::Message(
46 "prepared package index references a non-package remote object".to_string(),
47 ));
48 }
49 let remote_object_id = remote.object_id();
50 let object = remote.object().clone();
51 let coven_protocol::remote_object::SemanticPayload::Spooled(semantic_hash) =
52 remote.semantic_payload()
53 else {
54 return Err(DbError::Message(
55 "prepared package remote object names no stored plaintext".to_string(),
56 ));
57 };
58 let stored_hash = remote.stored_payload().ok_or_else(|| {
59 DbError::Message("prepared package remote object uploads no ciphertext".to_string())
60 })?;
61 let read = |hash| {
62 crate::payload_store::read_payload_blocking(conn, store_dir, hash)
63 .map_err(DbError::from)
64 };
65 let semantic_bytes = read(semantic_hash)?;
66 let stored_bytes = read(stored_hash)?;
67 Self::new(remote_object_id, semantic_bytes, stored_bytes, object)
68 }
69
70 pub fn new(
71 remote_object_id: ObjectHash,
72 semantic_bytes: Vec<u8>,
73 stored_bytes: Vec<u8>,
74 object: ExactObjectRef,
75 ) -> Result<Self, DbError> {
76 let package = AudiencePackage::parse(&semantic_bytes)
77 .map_err(|error| DbError::context("prepared audience package", error))?;
78 object
79 .verify(&stored_bytes)
80 .map_err(|error| DbError::context("prepared audience package stored bytes", error))?;
81 Ok(Self {
82 remote_object_id,
83 package,
84 semantic_bytes,
85 stored_bytes,
86 object,
87 })
88 }
89
90 pub fn remote_object_id(&self) -> ObjectHash {
91 self.remote_object_id
92 }
93
94 pub fn package(&self) -> &AudiencePackage {
95 &self.package
96 }
97
98 pub fn semantic_bytes(&self) -> &[u8] {
99 &self.semantic_bytes
100 }
101
102 pub fn stored_bytes(&self) -> &[u8] {
103 &self.stored_bytes
104 }
105
106 pub fn object(&self) -> &ExactObjectRef {
107 &self.object
108 }
109}
110
111#[derive(Debug, Clone, PartialEq, Eq)]
112pub struct PreparedAudienceBlob {
113 remote_object_id: ObjectHash,
114 audience: RemoteAudience,
115 blob: StoredBlobRef,
116 spool_path: Option<PathBuf>,
117}
118
119impl PreparedAudienceBlob {
120 pub fn from_remote(
121 audience: RemoteAudience,
122 expected_locator_hash: &str,
123 remote: RemoteObjectRecord,
124 spool_path: Option<PathBuf>,
125 ) -> Result<Self, DbError> {
126 remote
127 .validate()
128 .map_err(|error| DbError::context("prepared remote blob", error))?;
129 if !matches!(
130 &remote,
131 RemoteObjectRecord::SharedLiveSet(record)
132 if record.identity.domain == SharedLiveSetObjectDomain::StoredBlob
133 ) {
134 return Err(DbError::Message(
135 "prepared blob index references a non-blob remote object".to_string(),
136 ));
137 }
138 let locator_bytes = remote.payloads().carried_locator_bytes().ok_or_else(|| {
139 DbError::Message("prepared blob remote object carries no locator".to_string())
140 })?;
141 let locator = BlobLocator::parse(locator_bytes)
142 .map_err(|error| DbError::context("prepared blob locator", error))?;
143 if locator.locator_hash().to_string() != expected_locator_hash {
144 return Err(DbError::Message(format!(
145 "prepared blob locator hashes to {}, indexed as {expected_locator_hash}",
146 locator.locator_hash()
147 )));
148 }
149 if locator.audience() != audience {
150 return Err(DbError::Message(format!(
151 "prepared blob index audience {audience:?} differs from locator audience {:?}",
152 locator.audience()
153 )));
154 }
155 let requires_upload = matches!(
156 &remote,
157 RemoteObjectRecord::SharedLiveSet(record)
158 if matches!(record.state, coven_protocol::remote_object::OwnedObjectState::Prepared { .. })
159 );
160 if requires_upload && spool_path.is_none() {
161 return Err(DbError::Message(
162 "prepared blob awaiting upload has no local spool".to_string(),
163 ));
164 }
165 if spool_path.as_ref().is_some_and(|path| !path.is_absolute()) {
166 return Err(DbError::Message(
167 "prepared blob local spool path is not absolute".to_string(),
168 ));
169 }
170 let blob = StoredBlobRef::new(locator, remote.object().clone())
171 .map_err(|error| DbError::context("prepared blob reference", error))?;
172 Ok(Self {
173 remote_object_id: remote.object_id(),
174 audience,
175 blob,
176 spool_path,
177 })
178 }
179
180 pub fn remote_object_id(&self) -> ObjectHash {
181 self.remote_object_id
182 }
183
184 pub fn audience(&self) -> &RemoteAudience {
185 &self.audience
186 }
187
188 pub fn blob(&self) -> &StoredBlobRef {
189 &self.blob
190 }
191
192 pub fn spool_path(&self) -> Option<&Path> {
193 self.spool_path.as_deref()
194 }
195}
196
197#[derive(Debug, Clone)]
198pub struct PreparedAudienceObjects {
199 pub packages: Vec<PreparedAudiencePackage>,
200 pub blobs: Vec<PreparedAudienceBlob>,
201}
202
203pub struct PreparedRemoteObject {
204 pub closed: coven_protocol::remote_object::ClosedRemoteObject,
207 pub spool_path: Option<PathBuf>,
208}
209
210#[derive(Debug, Clone, PartialEq, Eq)]
211pub enum MakeRemoteIntentState {
212 Uploading,
213 Cancelling,
214 Publishing(WriteId),
215}
216
217#[derive(Debug, Clone, Copy, PartialEq, Eq)]
218pub enum StoredBlobReferenceState {
219 NotLiveRemote,
220 LiveRemote,
221 Unresolved,
222}
223
224pub fn validate_prepared_audience_blob_graph(
225 object_ids: &std::collections::BTreeSet<ObjectHash>,
226 audiences: &PreparedAudienceObjects,
227) -> Result<(), DbError> {
228 let mut indexed = std::collections::BTreeSet::new();
229 for package in &audiences.packages {
230 if !indexed.insert(package.remote_object_id()) {
231 return Err(DbError::Message(
232 "prepared audience objects contain a duplicate package body".to_string(),
233 ));
234 }
235 }
236 for blob in &audiences.blobs {
237 if !indexed.insert(blob.remote_object_id()) {
238 return Err(DbError::Message(
239 "prepared audience objects contain a duplicate blob body".to_string(),
240 ));
241 }
242 }
243 if &indexed != object_ids {
244 return Err(DbError::Message(
245 "closed remote objects differ from package/blob indexes".to_string(),
246 ));
247 }
248 validate_prepared_audience_blob_bindings(audiences)
249}
250
251pub(crate) fn validate_prepared_audience_blob_bindings(
252 audiences: &PreparedAudienceObjects,
253) -> Result<(), DbError> {
254 for package in &audiences.packages {
255 let audience = package.package().audience().remote_audience();
256 for binding in package.package().blob_bindings() {
257 if !audiences
258 .blobs
259 .iter()
260 .any(|blob| blob.audience() == &audience && blob.blob() == binding.blob())
261 {
262 return Err(DbError::Message(
263 "prepared package blob binding has no exact blob index".to_string(),
264 ));
265 }
266 }
267 }
268 for blob in &audiences.blobs {
269 if !audiences.packages.iter().any(|package| {
270 package.package().audience().remote_audience() == *blob.audience()
271 && package
272 .package()
273 .blob_bindings()
274 .iter()
275 .any(|binding| binding.blob() == blob.blob())
276 }) {
277 return Err(DbError::Message(
278 "prepared blob index has no exact package binding".to_string(),
279 ));
280 }
281 }
282 Ok(())
283}
284
285#[cfg(test)]
286mod tests {
287 use super::*;
288
289 fn exercise_exact_outbound_blob_graph(
292 circle: bool,
293 include_body: bool,
294 include_locator: bool,
295 include_binding: bool,
296 ) -> Result<(), DbError> {
297 use coven_protocol::audience_package::RowBlobLocatorBinding;
298 use coven_protocol::blob::BlobScope;
299 use coven_protocol::causal_grants::AuthorStreamId;
300 use coven_protocol::circle::CircleId;
301 use coven_protocol::circle_control::CircleControlCoord;
302 use coven_protocol::objects::ObjectSlot;
303 use coven_protocol::store_commit::{CandidateFamilyId, StoreCommitCoord};
304
305 let store_root_hash = ObjectHash::digest(b"outbound-graph-store");
306 let write_id = WriteId::from_generated("outbound-graph-write".to_string());
307 let coord = StoreCommitCoord {
308 stream_id: AuthorStreamId::from_bytes([4; 32]),
309 sequence: 1,
310 };
311 let candidate_family = CandidateFamilyId::from_hash(ObjectHash::digest(b"outbound-family"));
312 let remote_audience = if circle {
313 RemoteAudience::Circle(CircleId::from_bytes([7; 16]))
314 } else {
315 RemoteAudience::Store
316 };
317 let uploader_bytes = b"outbound graph uploader registration";
318 let uploader = StoreDeviceRegistrationRef {
319 device_id: "01"
320 .repeat(32)
321 .parse::<coven_protocol::store_commit::StoreDeviceId>()?,
322 registration_hash: ObjectHash::digest(uploader_bytes),
323 object: ExactObjectRef::new(
324 ObjectSlot::logical("store-v1/registrations/outbound-graph.json".to_string())?,
325 uploader_bytes.len() as u64,
326 ObjectHash::digest(uploader_bytes),
327 ),
328 };
329 let locator = BlobLocator::opaque(
330 "media".to_string(),
331 "blob-a".to_string(),
332 uploader,
333 remote_audience.clone(),
334 BlobScope::Master,
335 coven_keys::encryption::KeyFingerprint::from_bytes([3; 32]),
336 7,
337 ObjectHash::digest(b"content"),
338 )?;
339 let stored_bytes = b"sealed-content".to_vec();
340 let object = ExactObjectRef::new(
341 ObjectSlot::logical(locator.semantic_key())?,
342 stored_bytes.len() as u64,
343 ObjectHash::digest(&stored_bytes),
344 );
345 let stored = StoredBlobRef::new(locator, object)?;
346 let bindings = if include_binding {
347 vec![RowBlobLocatorBinding::new(
348 "items",
349 "row-a",
350 "stamp-a",
351 "media_blob",
352 stored.clone(),
353 )?]
354 } else {
355 Vec::new()
356 };
357 let package = if let RemoteAudience::Circle(circle_id) = remote_audience {
358 AudiencePackage::circle(
359 store_root_hash,
360 candidate_family,
361 write_id.clone(),
362 coord,
363 1,
364 circle_id,
365 CircleControlCoord {
366 device_id: "01".repeat(32),
367 stream_id: AuthorStreamId::from_bytes([5; 32]),
368 author_pubkey: "author-a".to_string(),
369 author_owner_grant: coven_protocol::causal_grants::MembershipGrantId(
370 ObjectHash::digest(b"outbound-graph owner grant"),
371 ),
372 seq: 1,
373 control_hash: ObjectHash::digest(b"circle-control"),
374 },
375 coven_keys::encryption::KeyFingerprint::from_bytes([3; 32]),
376 b"changeset".to_vec(),
377 bindings,
378 )
379 } else {
380 AudiencePackage::store(
381 store_root_hash,
382 candidate_family,
383 write_id,
384 coord,
385 1,
386 b"changeset".to_vec(),
387 bindings,
388 )
389 }?;
390 let package_bytes = package.to_bytes();
391 let package_object = ExactObjectRef::new(
392 ObjectSlot::logical("test/package".to_string())?,
393 package_bytes.len() as u64,
394 ObjectHash::digest(&package_bytes),
395 );
396 let package_id = ObjectHash::digest(b"package-record");
397 let blob_id = ObjectHash::digest(b"blob-record");
398 let packages = vec![PreparedAudiencePackage::new(
399 package_id,
400 package_bytes.clone(),
401 package_bytes,
402 package_object,
403 )?];
404 let blobs = if include_locator {
405 vec![PreparedAudienceBlob {
406 remote_object_id: blob_id,
407 audience: remote_audience,
408 blob: stored,
409 spool_path: Some(PathBuf::from("/outbound-blob.spool")),
410 }]
411 } else {
412 Vec::new()
413 };
414 let mut object_ids = std::collections::BTreeSet::from([package_id]);
415 if include_body {
416 object_ids.insert(blob_id);
417 }
418 validate_prepared_audience_blob_graph(
419 &object_ids,
420 &PreparedAudienceObjects { packages, blobs },
421 )
422 }
423
424 #[test]
429 fn store_and_circle_blob_publication_require_body_locator_and_binding() {
430 for circle in [false, true] {
431 assert!(exercise_exact_outbound_blob_graph(circle, false, true, true).is_err());
432 assert!(exercise_exact_outbound_blob_graph(circle, true, false, true).is_err());
433 assert!(exercise_exact_outbound_blob_graph(circle, true, true, false).is_err());
434 exercise_exact_outbound_blob_graph(circle, true, true, true).unwrap();
435 }
436 }
437}