coven_replication/sync/store/circles/
packages.rs1use 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}