Skip to main content

coven_replication/sync/store/circles/
packages.rs

1use tracing::debug;
2
3use crate::sync::store::commit_verification::merge_history::MergeHistoryVerifier;
4use crate::sync::store::pull::{LoadedCirclePackage, LocalStoreMembership};
5use coven_database::store::CirclePackageAccess;
6use coven_database::{DbError, StoreDatabase};
7use coven_protocol::circle_activation::{
8    VerifiedCircleActivations, VerifiedStreamActivationPrefix,
9};
10use coven_protocol::objects::VerifiedObject;
11use coven_protocol::store_commit::{
12    CirclePackageRef, StoreDeviceRegistration, StoreProtocolError, VerifiedStoreBatchCommit,
13};
14use coven_storage::run_blocking_object_verification;
15
16#[derive(Debug, thiserror::Error)]
17pub enum CirclePackageReadError {
18    #[error("Circle package database: {0}")]
19    Database(#[from] DbError),
20    #[error("Circle package is invalid: {0}")]
21    Invalid(String),
22    #[error("Circle package state: {0}")]
23    CircleState(#[from] coven_protocol::circle_activation::CircleStateError),
24    #[error("Circle package storage: {0}")]
25    Storage(#[from] coven_protocol::objects::StorageError),
26    #[error("Circle package object: {0}")]
27    StoreObject(#[from] coven_protocol::objects::StoreObjectError),
28    #[error("Circle package sync cycle: {0}")]
29    SyncCycle(#[source] Box<crate::sync::cycle::SyncCycleFailure>),
30    #[error("Circle package operation: {0}")]
31    CircleOperation(#[source] Box<crate::sync::store::circles::CircleOperationError>),
32    #[error("Circle package pull: {0}")]
33    Pull(#[source] Box<crate::sync::store::StorePullError>),
34    #[error("Circle package roster: {0}")]
35    Roster(#[from] coven_protocol::circle_roster::CircleRosterError),
36}
37
38impl From<crate::sync::cycle::SyncCycleFailure> for CirclePackageReadError {
39    fn from(error: crate::sync::cycle::SyncCycleFailure) -> Self {
40        Self::SyncCycle(Box::new(error))
41    }
42}
43
44impl From<crate::sync::store::circles::CircleOperationError> for CirclePackageReadError {
45    fn from(error: crate::sync::store::circles::CircleOperationError) -> Self {
46        Self::CircleOperation(Box::new(error))
47    }
48}
49
50impl From<crate::sync::store::StorePullError> for CirclePackageReadError {
51    fn from(error: crate::sync::store::StorePullError) -> Self {
52        Self::Pull(Box::new(error))
53    }
54}
55
56pub(crate) struct OpenedCirclePackage {
57    pub(crate) object: VerifiedObject<Vec<u8>>,
58}
59
60pub(crate) struct CirclePackageReader<'operation, 'storage> {
61    database: &'operation StoreDatabase,
62    storage: &'storage dyn coven_storage::CloudSyncObjectStorage,
63    history: &'operation mut MergeHistoryVerifier<'storage>,
64}
65
66impl<'operation, 'storage> CirclePackageReader<'operation, 'storage> {
67    pub(crate) fn new(
68        database: &'operation StoreDatabase,
69        storage: &'storage dyn coven_storage::CloudSyncObjectStorage,
70        history: &'operation mut MergeHistoryVerifier<'storage>,
71    ) -> Self {
72        Self {
73            database,
74            storage,
75            history,
76        }
77    }
78
79    fn root(&self) -> &coven_protocol::store_commit::StoreRootRef {
80        self.history.verified_root().reference()
81    }
82
83    pub(crate) async fn open_package(
84        &self,
85        access: &coven_protocol::circle_activation::CircleEpochAccess,
86        verified: &VerifiedStoreBatchCommit,
87        reference: &CirclePackageRef,
88        author: &StoreDeviceRegistration,
89    ) -> Result<OpenedCirclePackage, CirclePackageReadError> {
90        access
91            .authorize_package(reference, author)
92            .map_err(CirclePackageReadError::from)?;
93        let commit = verified.value();
94        if !commit
95            .circle_packages()
96            .iter()
97            .any(|committed| committed == reference)
98        {
99            return Err(CirclePackageReadError::Invalid(
100                StoreProtocolError::MissingCirclePackage(reference.circle_id).to_string(),
101            ));
102        }
103        let semantic_prefix = coven_protocol::store_commit::circle_package_semantic_prefix(
104            reference.circle_id,
105            commit.candidate_family(),
106            &verified.reference().coord.stream_id.to_string(),
107            commit.seq(),
108            reference.package.content_hash,
109        );
110        let context = access.protocol_context(
111            commit.store_root_hash,
112            coven_protocol::objects::ProtocolObjectDomain::CirclePackage,
113        );
114        let bytes = self
115            .storage
116            .read_protocol_object(&context, &reference.package.object, &semantic_prefix)
117            .await
118            .map_err(CirclePackageReadError::from)?;
119        let verify_bytes = bytes.clone();
120        let expected_commit = commit.clone();
121        let expected_circle_id = reference.circle_id;
122        let value = run_blocking_object_verification(
123            &semantic_prefix,
124            &reference.package.object,
125            Box::new(move || {
126                expected_commit.verify_circle_package(expected_circle_id, &verify_bytes)?;
127                Ok(verify_bytes)
128            }),
129        )
130        .await
131        .map_err(CirclePackageReadError::from)?;
132        Ok(OpenedCirclePackage {
133            object: VerifiedObject {
134                value,
135                bytes,
136                semantic_hash: reference.package.content_hash,
137                object: reference.package.object.clone(),
138            },
139        })
140    }
141
142    pub(crate) async fn load_applicable(
143        &mut self,
144        verified: &VerifiedStoreBatchCommit,
145        prepared: &[&VerifiedCircleActivations],
146        snapshot_access: &[coven_database::StagedCircleAccess],
147        author: &StoreDeviceRegistration,
148        local_store_membership: LocalStoreMembership,
149    ) -> Result<Vec<LoadedCirclePackage>, CirclePackageReadError> {
150        self.load_selected(
151            verified,
152            verified.value().circle_packages(),
153            prepared,
154            snapshot_access,
155            author,
156            local_store_membership,
157        )
158        .await
159    }
160
161    pub(crate) async fn load_selected(
162        &mut self,
163        verified: &VerifiedStoreBatchCommit,
164        references: &[CirclePackageRef],
165        prepared: &[&VerifiedCircleActivations],
166        snapshot_access: &[coven_database::StagedCircleAccess],
167        author: &StoreDeviceRegistration,
168        local_store_membership: LocalStoreMembership,
169    ) -> Result<Vec<LoadedCirclePackage>, CirclePackageReadError> {
170        let root = self.root().clone();
171        let commit_ref = verified.reference();
172        let commit = verified.value();
173        if references.is_empty() {
174            return Ok(Vec::new());
175        }
176        let activations = prepared
177            .iter()
178            .flat_map(|group| group.circles())
179            .chain(snapshot_access.iter().map(|access| &access.activation))
180            .cloned()
181            .collect::<Vec<_>>();
182        let mut verified_prefix = VerifiedStreamActivationPrefix::empty();
183        for group in prepared {
184            verified_prefix
185                .include(group.stream_activations())
186                .map_err(|error| CirclePackageReadError::Invalid(error.to_string()))?;
187        }
188        let mut replay_epochs = self
189            .database
190            .circle_replay_epoch_index(root.clone())
191            .await
192            .map_err(CirclePackageReadError::Database)?;
193        replay_epochs
194            .include_verified_activations(&activations)
195            .map_err(CirclePackageReadError::Database)?;
196        let mut loaded = Vec::new();
197        for reference in references {
198            if !replay_epochs
199                .permits(commit_ref, reference.circle_id, &reference.control)
200                .map_err(CirclePackageReadError::Database)?
201            {
202                debug!(
203                    circle_id = %reference.circle_id,
204                    control = ?reference.control,
205                    "skipping Circle package beyond its accepted epoch cutoff"
206                );
207                continue;
208            }
209            if matches!(
210                local_store_membership,
211                LocalStoreMembership::IdentityNotSupplied
212            ) {
213                return Err(CirclePackageReadError::Invalid(format!(
214                    "commit {} carries Circle packages but no verified local Store membership was supplied",
215                    commit.seq()
216                )));
217            }
218            if matches!(local_store_membership, LocalStoreMembership::Removed) {
219                debug!(
220                    circle_id = %reference.circle_id,
221                    control = ?reference.control,
222                    "skipping Circle package for an identity removed from Store membership"
223                );
224                continue;
225            }
226            if reference.package.schema_version > self.database.schema_version() {
227                return Err(CirclePackageReadError::Invalid(format!(
228                    "Circle package for {} requires schema {}, local schema is {}",
229                    reference.circle_id,
230                    reference.package.schema_version,
231                    self.database.schema_version()
232                )));
233            }
234            let Some(package_access) = self
235                .database
236                .circle_package_access(
237                    root.clone(),
238                    reference.circle_id,
239                    reference.control.clone(),
240                    reference.key_fingerprint,
241                    activations.clone(),
242                )
243                .await
244                .map_err(CirclePackageReadError::Database)?
245            else {
246                debug!(circle_id = %reference.circle_id, control = ?reference.control,
247                    "skipping Circle package without applicable local access");
248                continue;
249            };
250            let access = match package_access {
251                CirclePackageAccess::Exact(access) => access,
252                CirclePackageAccess::Historical(keyring) => {
253                    let prepared_historical = prepared
254                        .iter()
255                        .find_map(|group| {
256                            group
257                                .circles()
258                                .iter()
259                                .find(|activation| {
260                                    activation.circle_id == reference.circle_id
261                                        && activation.control.coord == reference.control
262                                })
263                                .map(|activation| {
264                                    (
265                                        activation.clone(),
266                                        group.stream_activations().activating_commit().clone(),
267                                    )
268                                })
269                        })
270                        .or_else(|| {
271                            snapshot_access
272                                .iter()
273                                .find(|access| {
274                                    access.activation.circle_id == reference.circle_id
275                                        && access.activation.control.coord == reference.control
276                                })
277                                .map(|access| {
278                                    (access.activation.clone(), access.activating_commit.clone())
279                                })
280                        });
281                    let historical_context = match prepared_historical {
282                        Some(context) => Some(context),
283                        None => self
284                            .database
285                            .verified_circle_activation_context(
286                                root.clone(),
287                                reference.circle_id,
288                                reference.control.clone(),
289                            )
290                            .await
291                            .map_err(CirclePackageReadError::Database)?,
292                    };
293                    let Some((historical, historical_commit_ref)) = historical_context else {
294                        return Err(CirclePackageReadError::Invalid(format!(
295                            "Circle {} historical package control is not retained",
296                            reference.circle_id
297                        )));
298                    };
299                    let historical_commit = self.history.load_ref(&historical_commit_ref).await?;
300                    let roster_chain = super::activation::CircleActivationVerifier::new(
301                        self.database,
302                        self.storage,
303                        self.history,
304                    )
305                    .load_control_roster_chain_with_prefix(
306                        &historical_commit,
307                        &historical.reference,
308                        &historical.control,
309                        &keyring,
310                        &verified_prefix,
311                    )
312                    .await
313                    .map_err(CirclePackageReadError::from)?;
314                    let roster = roster_chain.try_resolved()?;
315                    coven_protocol::circle_activation::CircleEpochAccess::from_historical(
316                        reference.circle_id,
317                        reference.key_fingerprint,
318                        &keyring,
319                        &roster,
320                    )
321                    .map_err(CirclePackageReadError::from)?
322                }
323            };
324            let package = self
325                .open_package(&access, verified, reference, author)
326                .await?;
327            loaded.push(LoadedCirclePackage {
328                reference: reference.clone(),
329                bytes: package.object.value,
330            });
331        }
332        Ok(loaded)
333    }
334}