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::{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}