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