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, serde::Serialize, serde::Deserialize)]
112#[serde(deny_unknown_fields)]
113pub struct PreparedAudienceBlob {
114 remote_object_id: ObjectHash,
115 audience: RemoteAudience,
116 blob: StoredBlobRef,
117 spool_path: Option<PathBuf>,
118}
119
120impl PreparedAudienceBlob {
121 pub fn from_remote(
122 audience: RemoteAudience,
123 expected_locator_hash: &str,
124 remote: RemoteObjectRecord,
125 spool_path: Option<PathBuf>,
126 ) -> Result<Self, DbError> {
127 remote
128 .validate()
129 .map_err(|error| DbError::context("prepared remote blob", error))?;
130 if !matches!(
131 &remote,
132 RemoteObjectRecord::SharedLiveSet(record)
133 if record.identity.domain == SharedLiveSetObjectDomain::StoredBlob
134 ) {
135 return Err(DbError::Message(
136 "prepared blob index references a non-blob remote object".to_string(),
137 ));
138 }
139 let locator_bytes = remote.payloads().carried_locator_bytes().ok_or_else(|| {
140 DbError::Message("prepared blob remote object carries no locator".to_string())
141 })?;
142 let locator = BlobLocator::parse(locator_bytes)
143 .map_err(|error| DbError::context("prepared blob locator", error))?;
144 if locator.locator_hash().to_string() != expected_locator_hash {
145 return Err(DbError::Message(format!(
146 "prepared blob locator hashes to {}, indexed as {expected_locator_hash}",
147 locator.locator_hash()
148 )));
149 }
150 if locator.audience() != audience {
151 return Err(DbError::Message(format!(
152 "prepared blob index audience {audience:?} differs from locator audience {:?}",
153 locator.audience()
154 )));
155 }
156 let requires_upload = matches!(
157 &remote,
158 RemoteObjectRecord::SharedLiveSet(record)
159 if matches!(record.state, coven_protocol::remote_object::OwnedObjectState::Prepared { .. })
160 );
161 if requires_upload && spool_path.is_none() {
162 return Err(DbError::Message(
163 "prepared blob awaiting upload has no local spool".to_string(),
164 ));
165 }
166 if spool_path.as_ref().is_some_and(|path| !path.is_absolute()) {
167 return Err(DbError::Message(
168 "prepared blob local spool path is not absolute".to_string(),
169 ));
170 }
171 let blob = StoredBlobRef::new(locator, remote.object().clone())
172 .map_err(|error| DbError::context("prepared blob reference", error))?;
173 Ok(Self {
174 remote_object_id: remote.object_id(),
175 audience,
176 blob,
177 spool_path,
178 })
179 }
180
181 pub fn remote_object_id(&self) -> ObjectHash {
182 self.remote_object_id
183 }
184
185 pub fn audience(&self) -> &RemoteAudience {
186 &self.audience
187 }
188
189 pub fn blob(&self) -> &StoredBlobRef {
190 &self.blob
191 }
192
193 pub fn spool_path(&self) -> Option<&Path> {
194 self.spool_path.as_deref()
195 }
196}
197
198#[derive(Debug, Clone)]
199pub struct PreparedAudienceObjects {
200 pub packages: Vec<PreparedAudiencePackage>,
201 pub blobs: Vec<PreparedAudienceBlob>,
202}
203
204pub struct PreparedRemoteObject {
205 pub closed: coven_protocol::remote_object::ClosedRemoteObject,
208 pub spool_path: Option<PathBuf>,
209}
210
211#[derive(Debug, Clone, PartialEq, Eq)]
212pub enum MakeRemoteIntentState {
213 Uploading,
214 Cancelling,
215 Publishing(WriteId),
216}
217
218#[derive(Debug, Clone, Copy, PartialEq, Eq)]
219pub enum StoredBlobReferenceState {
220 NotLiveRemote,
221 LiveRemote,
222 Unresolved,
223}
224
225pub fn validate_prepared_audience_blob_graph(
226 object_ids: &std::collections::BTreeSet<ObjectHash>,
227 audiences: &PreparedAudienceObjects,
228) -> Result<(), DbError> {
229 let mut indexed = std::collections::BTreeSet::new();
230 for package in &audiences.packages {
231 if !indexed.insert(package.remote_object_id()) {
232 return Err(DbError::Message(
233 "prepared audience objects contain a duplicate package body".to_string(),
234 ));
235 }
236 }
237 for blob in &audiences.blobs {
238 if !indexed.insert(blob.remote_object_id()) {
239 return Err(DbError::Message(
240 "prepared audience objects contain a duplicate blob body".to_string(),
241 ));
242 }
243 }
244 if &indexed != object_ids {
245 return Err(DbError::Message(
246 "closed remote objects differ from package/blob indexes".to_string(),
247 ));
248 }
249 validate_prepared_audience_blob_bindings(audiences)
250}
251
252pub(crate) fn validate_prepared_audience_blob_bindings(
253 audiences: &PreparedAudienceObjects,
254) -> Result<(), DbError> {
255 for package in &audiences.packages {
256 let audience = package.package().audience().remote_audience();
257 for binding in package.package().blob_bindings() {
258 if !audiences
259 .blobs
260 .iter()
261 .any(|blob| blob.audience() == &audience && blob.blob() == binding.blob())
262 {
263 return Err(DbError::Message(
264 "prepared package blob binding has no exact blob index".to_string(),
265 ));
266 }
267 }
268 }
269 for blob in &audiences.blobs {
270 if !audiences.packages.iter().any(|package| {
271 package.package().audience().remote_audience() == *blob.audience()
272 && package
273 .package()
274 .blob_bindings()
275 .iter()
276 .any(|binding| binding.blob() == blob.blob())
277 }) {
278 return Err(DbError::Message(
279 "prepared blob index has no exact package binding".to_string(),
280 ));
281 }
282 }
283 Ok(())
284}
285
286#[cfg(test)]
287mod tests {
288 use super::*;
289
290 fn exercise_exact_outbound_blob_graph(
293 circle: bool,
294 include_body: bool,
295 include_locator: bool,
296 include_binding: bool,
297 ) -> Result<(), DbError> {
298 use coven_protocol::audience_package::RowBlobLocatorBinding;
299 use coven_protocol::blob::BlobScope;
300 use coven_protocol::causal_grants::AuthorStreamId;
301 use coven_protocol::circle::CircleId;
302 use coven_protocol::circle_control::CircleControlCoord;
303 use coven_protocol::objects::ObjectSlot;
304 use coven_protocol::store_commit::{CandidateFamilyId, StoreCommitCoord};
305
306 let store_root_hash = ObjectHash::digest(b"outbound-graph-store");
307 let write_id = WriteId::from_generated("outbound-graph-write".to_string());
308 let coord = StoreCommitCoord {
309 stream_id: AuthorStreamId::from_bytes([4; 32]),
310 sequence: 1,
311 };
312 let candidate_family = CandidateFamilyId::from_hash(ObjectHash::digest(b"outbound-family"));
313 let remote_audience = if circle {
314 RemoteAudience::Circle(CircleId::from_bytes([7; 16]))
315 } else {
316 RemoteAudience::Store
317 };
318 let uploader_bytes = b"outbound graph uploader registration";
319 let uploader = StoreDeviceRegistrationRef {
320 device_id: "01"
321 .repeat(32)
322 .parse::<coven_protocol::store_commit::StoreDeviceId>()?,
323 registration_hash: ObjectHash::digest(uploader_bytes),
324 object: ExactObjectRef::new(
325 ObjectSlot::logical("store-v1/registrations/outbound-graph.json".to_string())?,
326 uploader_bytes.len() as u64,
327 ObjectHash::digest(uploader_bytes),
328 ),
329 };
330 let locator = BlobLocator::opaque(
331 "media".to_string(),
332 "blob-a".to_string(),
333 uploader,
334 remote_audience.clone(),
335 BlobScope::Master,
336 coven_keys::encryption::KeyFingerprint::from_bytes([3; 32]),
337 7,
338 ObjectHash::digest(b"content"),
339 )?;
340 let stored_bytes = b"sealed-content".to_vec();
341 let object = ExactObjectRef::new(
342 ObjectSlot::logical(locator.semantic_key())?,
343 stored_bytes.len() as u64,
344 ObjectHash::digest(&stored_bytes),
345 );
346 let stored = StoredBlobRef::new(locator, object)?;
347 let bindings = if include_binding {
348 vec![RowBlobLocatorBinding::new(
349 "items",
350 "row-a",
351 "stamp-a",
352 "media_blob",
353 stored.clone(),
354 )?]
355 } else {
356 Vec::new()
357 };
358 let package = if let RemoteAudience::Circle(circle_id) = remote_audience {
359 AudiencePackage::circle(
360 store_root_hash,
361 candidate_family,
362 write_id.clone(),
363 coord,
364 1,
365 circle_id,
366 CircleControlCoord {
367 device_id: "01".repeat(32),
368 stream_id: AuthorStreamId::from_bytes([5; 32]),
369 author_pubkey: "author-a".to_string(),
370 author_owner_grant: coven_protocol::causal_grants::MembershipGrantId(
371 ObjectHash::digest(b"outbound-graph owner grant"),
372 ),
373 seq: 1,
374 control_hash: ObjectHash::digest(b"circle-control"),
375 },
376 coven_keys::encryption::KeyFingerprint::from_bytes([3; 32]),
377 b"changeset".to_vec(),
378 bindings,
379 )
380 } else {
381 AudiencePackage::store(
382 store_root_hash,
383 candidate_family,
384 write_id,
385 coord,
386 1,
387 b"changeset".to_vec(),
388 bindings,
389 )
390 }?;
391 let package_bytes = package.to_bytes();
392 let package_object = ExactObjectRef::new(
393 ObjectSlot::logical("test/package".to_string())?,
394 package_bytes.len() as u64,
395 ObjectHash::digest(&package_bytes),
396 );
397 let package_id = ObjectHash::digest(b"package-record");
398 let blob_id = ObjectHash::digest(b"blob-record");
399 let packages = vec![PreparedAudiencePackage::new(
400 package_id,
401 package_bytes.clone(),
402 package_bytes,
403 package_object,
404 )?];
405 let blobs = if include_locator {
406 vec![PreparedAudienceBlob {
407 remote_object_id: blob_id,
408 audience: remote_audience,
409 blob: stored,
410 spool_path: Some(PathBuf::from("/outbound-blob.spool")),
411 }]
412 } else {
413 Vec::new()
414 };
415 let mut object_ids = std::collections::BTreeSet::from([package_id]);
416 if include_body {
417 object_ids.insert(blob_id);
418 }
419 validate_prepared_audience_blob_graph(
420 &object_ids,
421 &PreparedAudienceObjects { packages, blobs },
422 )
423 }
424
425 #[test]
430 fn store_and_circle_blob_publication_require_body_locator_and_binding() {
431 for circle in [false, true] {
432 assert!(exercise_exact_outbound_blob_graph(circle, false, true, true).is_err());
433 assert!(exercise_exact_outbound_blob_graph(circle, true, false, true).is_err());
434 assert!(exercise_exact_outbound_blob_graph(circle, true, true, false).is_err());
435 exercise_exact_outbound_blob_graph(circle, true, true, true).unwrap();
436 }
437 }
438}