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)]
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    /// The record awaiting upload, with the payloads its row names read
205    /// back beside it: the upload reads the ciphertext from here.
206    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    /// Build one prepared outbound graph and validate it, leaving out whichever
290    /// of the blob body, its locator, or its row binding the caller drops.
291    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    /// A publishable blob needs all three of its parts: the body among the
425    /// uploaded object ids, the locator in the prepared blobs, and a package
426    /// binding that names it. Any one missing is refused, for Store and Circle
427    /// audiences alike.
428    #[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}