Skip to main content

coven_database/
prepared_audience_objects.rs

1use super::*;
2
3/// The exact commit coordinate that first made a blob locator authoritative.
4#[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    /// The prepared package one remote object names, read back from the payload
20    /// spool the record's identity files it under.
21    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    /// The record awaiting upload, with the payloads its row names read
206    /// back beside it: the upload reads the ciphertext from here.
207    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    /// Build one prepared outbound graph and validate it, leaving out whichever
291    /// of the blob body, its locator, or its row binding the caller drops.
292    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    /// A publishable blob needs all three of its parts: the body among the
426    /// uploaded object ids, the locator in the prepared blobs, and a package
427    /// binding that names it. Any one missing is refused, for Store and Circle
428    /// audiences alike.
429    #[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}