1use super::nonactivation::*;
2use super::*;
3
4#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
5#[serde(rename_all = "snake_case", deny_unknown_fields)]
6pub enum OwnedObjectState {
7 Prepared {
8 ownership: PendingCandidateOwnership,
9 },
10 UploadedVerified {
11 ownership: SharedObjectOwnership,
12 },
13 RetirementPending {
14 former_candidates: Vec<CandidateNonactivation>,
15 },
16}
17
18impl OwnedObjectState {
19 pub(super) fn validate(&self) -> Result<(), RemoteObjectRecordError> {
20 match self {
21 Self::Prepared { ownership } => ownership.validate(),
22 Self::UploadedVerified { ownership } => ownership.validate(),
23 Self::RetirementPending { former_candidates } => {
24 validate_nonactivations(former_candidates)
25 }
26 }
27 }
28}
29
30pub(super) fn merge_store_commit_owner(state: &mut OwnedObjectState, owner: &StoreBatchCommitRef) {
31 match state {
32 OwnedObjectState::Prepared { ownership } => {
33 let mut pending = ownership.pending.clone();
34 pending.remove(owner);
35 *state = OwnedObjectState::UploadedVerified {
36 ownership: SharedObjectOwnership {
37 pending,
38 activated: BTreeSet::from([SharedObjectOwner::StoreCommit(owner.clone())]),
39 nonactivated: ownership.nonactivated.clone(),
40 },
41 };
42 }
43 OwnedObjectState::UploadedVerified { ownership } => {
44 ownership.pending.remove(owner);
45 ownership
46 .activated
47 .insert(SharedObjectOwner::StoreCommit(owner.clone()));
48 }
49 OwnedObjectState::RetirementPending { former_candidates } => {
50 *state = OwnedObjectState::UploadedVerified {
51 ownership: SharedObjectOwnership {
52 pending: BTreeSet::new(),
53 activated: BTreeSet::from([SharedObjectOwner::StoreCommit(owner.clone())]),
54 nonactivated: former_candidates.clone(),
55 },
56 };
57 }
58 }
59}
60
61pub(super) fn merge_shared_owner(
62 state: &mut OwnedObjectState,
63 owner: SharedObjectOwner,
64) -> Result<(), RemoteObjectRecordError> {
65 match state {
66 OwnedObjectState::UploadedVerified { ownership } => {
67 ownership.activated.insert(owner);
68 Ok(())
69 }
70 OwnedObjectState::RetirementPending { former_candidates } => {
71 *state = OwnedObjectState::UploadedVerified {
72 ownership: SharedObjectOwnership {
73 pending: BTreeSet::new(),
74 activated: BTreeSet::from([owner]),
75 nonactivated: former_candidates.clone(),
76 },
77 };
78 Ok(())
79 }
80 OwnedObjectState::Prepared { .. } => Err(RemoteObjectRecordError::InvalidActivation),
81 }
82}
83
84#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
85#[serde(deny_unknown_fields)]
86pub struct PendingCandidateOwnership {
87 pub pending: BTreeSet<StoreBatchCommitRef>,
88 pub nonactivated: Vec<CandidateNonactivation>,
89}
90
91impl PendingCandidateOwnership {
92 pub(super) fn validate(&self) -> Result<(), RemoteObjectRecordError> {
93 if self.pending.is_empty() {
94 return Err(RemoteObjectRecordError::EmptyPendingOwnership);
95 }
96 validate_owner_partition(&self.pending, std::iter::empty(), &self.nonactivated)
97 }
98}
99
100#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(deny_unknown_fields)]
102pub struct SharedObjectOwnership {
103 pub pending: BTreeSet<StoreBatchCommitRef>,
104 pub activated: BTreeSet<SharedObjectOwner>,
105 pub nonactivated: Vec<CandidateNonactivation>,
106}
107
108impl SharedObjectOwnership {
109 pub(super) fn validate(&self) -> Result<(), RemoteObjectRecordError> {
110 if self.pending.is_empty() && self.activated.is_empty() {
111 Err(RemoteObjectRecordError::EmptyOwnership)
112 } else {
113 let activated_commits = self.activated.iter().filter_map(|owner| match owner {
114 SharedObjectOwner::StoreCommit(commit) => Some(commit),
115 SharedObjectOwner::Snapshot(_) | SharedObjectOwner::RetainedReplay(_) => None,
116 });
117 validate_owner_partition(&self.pending, activated_commits, &self.nonactivated)
118 }
119 }
120}
121
122#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
123#[serde(deny_unknown_fields)]
124pub struct CandidateOwnership {
125 pub pending: BTreeSet<StoreBatchCommitRef>,
126 pub activated: BTreeSet<StoreBatchCommitRef>,
127 pub nonactivated: Vec<CandidateNonactivation>,
128}
129
130impl CandidateOwnership {
131 pub(super) fn validate(&self) -> Result<(), RemoteObjectRecordError> {
132 if self.pending.is_empty() && self.activated.is_empty() {
133 return Err(RemoteObjectRecordError::EmptyOwnership);
134 }
135 validate_owner_partition(&self.pending, self.activated.iter(), &self.nonactivated)
136 }
137}
138
139#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
140#[serde(rename_all = "snake_case", deny_unknown_fields)]
141pub enum SharedObjectOwner {
142 StoreCommit(StoreBatchCommitRef),
143 Snapshot(SnapshotObjectOwner),
144 RetainedReplay(RetainedReplayOwner),
145}
146
147#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
148#[serde(rename_all = "snake_case", deny_unknown_fields)]
149pub enum RetainedReplayOwner {
150 Commit {
151 commit: StoreBatchCommitRef,
152 input_hash: ObjectHash,
153 },
154}
155
156impl RetainedReplayOwner {
157 pub fn commit(&self) -> &StoreBatchCommitRef {
158 match self {
159 Self::Commit { commit, .. } => commit,
160 }
161 }
162}
163
164#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
165#[serde(rename_all = "snake_case", deny_unknown_fields)]
166pub enum SnapshotObjectOwner {
167 Store {
168 metadata_slot: crate::objects::ObjectSlot,
169 },
170 Circle {
171 activation: StreamActivationId,
172 generation: u64,
173 },
174}
175
176impl RemoteObjectRecord {
177 pub fn merge_blob_activation(
178 &mut self,
179 stored: &crate::blob::locator::StoredBlobRef,
180 owner: &StoreBatchCommitRef,
181 ) -> Result<(), RemoteObjectRecordError> {
182 let Self::SharedLiveSet(record) = self else {
183 return Err(RemoteObjectRecordError::DomainMismatch);
184 };
185 let locator_bytes = stored.locator().to_bytes();
186 if record.identity.domain != SharedLiveSetObjectDomain::StoredBlob
187 || record.identity.semantic_hash != ObjectHash::digest(&locator_bytes)
188 || record.identity.object != *stored.object()
189 || record.payloads.carried_locator_bytes() != Some(locator_bytes.as_slice())
190 {
191 return Err(RemoteObjectRecordError::StoredReferenceMismatch);
192 }
193 merge_store_commit_owner(&mut record.state, owner);
194 self.validate()
195 }
196
197 pub fn merge_package_activation(
198 &mut self,
199 domain: &SharedLiveSetObjectDomain,
200 package: &crate::audience_package::AudiencePackage,
201 owner: &StoreBatchCommitRef,
202 ) -> Result<(), RemoteObjectRecordError> {
203 if !matches!(
204 domain,
205 SharedLiveSetObjectDomain::StorePackage { .. }
206 | SharedLiveSetObjectDomain::CirclePackage { .. }
207 ) {
208 return Err(RemoteObjectRecordError::DomainMismatch);
209 }
210 let Self::SharedLiveSet(record) = self else {
211 return Err(RemoteObjectRecordError::DomainMismatch);
212 };
213 let canonical_semantic_bytes = package.to_bytes();
214 if &record.identity.domain != domain
215 || record.identity.semantic_hash != ObjectHash::digest(&canonical_semantic_bytes)
216 || record.identity.object != *domain.package_object()?
217 || matches!(record.payloads, RemoteObjectPayloads::RowBlob { .. })
218 {
219 return Err(RemoteObjectRecordError::StoredReferenceMismatch);
220 }
221 merge_store_commit_owner(&mut record.state, owner);
222 self.validate()
223 }
224
225 pub fn merge_retained_replay_owner(
226 &mut self,
227 owner: RetainedReplayOwner,
228 ) -> Result<(), RemoteObjectRecordError> {
229 let Self::SharedLiveSet(record) = self else {
230 return Err(RemoteObjectRecordError::DomainMismatch);
231 };
232 let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
233 return Err(RemoteObjectRecordError::InvalidActivation);
234 };
235 ownership
236 .activated
237 .insert(SharedObjectOwner::RetainedReplay(owner));
238 self.validate()
239 }
240
241 pub fn remove_all_retained_replay_owners(&mut self) -> Result<(), RemoteObjectRecordError> {
242 let Self::SharedLiveSet(record) = self else {
243 return Ok(());
244 };
245 let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
246 return Ok(());
247 };
248 ownership
249 .activated
250 .retain(|owner| !matches!(owner, SharedObjectOwner::RetainedReplay(_)));
251 self.validate()
252 }
253
254 pub fn remove_retained_replay_owner(
255 &mut self,
256 owner: &RetainedReplayOwner,
257 ) -> Result<(), RemoteObjectRecordError> {
258 let Self::SharedLiveSet(record) = self else {
259 return Err(RemoteObjectRecordError::DomainMismatch);
260 };
261 let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
262 return Err(RemoteObjectRecordError::InvalidActivation);
263 };
264 if !ownership
265 .activated
266 .remove(&SharedObjectOwner::RetainedReplay(owner.clone()))
267 {
268 return Err(RemoteObjectRecordError::CandidateOwnerMismatch);
269 }
270 self.retire_unowned_shared_live_set()?;
271 self.validate()
272 }
273
274 fn retire_unowned_shared_live_set(&mut self) -> Result<(), RemoteObjectRecordError> {
275 let Self::SharedLiveSet(record) = self else {
276 return Err(RemoteObjectRecordError::DomainMismatch);
277 };
278 let OwnedObjectState::UploadedVerified { ownership } = &record.state else {
279 return Err(RemoteObjectRecordError::InvalidActivation);
280 };
281 if !ownership.pending.is_empty() || !ownership.activated.is_empty() {
282 return Ok(());
283 }
284 if ownership.nonactivated.is_empty() {
285 return Err(RemoteObjectRecordError::EmptyOwnership);
286 }
287 let former_candidates = ownership.nonactivated.clone();
288 let package_domain = match &record.identity.domain {
289 SharedLiveSetObjectDomain::StorePackage { reference } => Some((
290 reference.candidate_family,
291 CandidateExclusiveObjectDomain::StorePackage {
292 reference: reference.clone(),
293 },
294 )),
295 SharedLiveSetObjectDomain::CirclePackage { reference } => Some((
296 reference.package.candidate_family,
297 CandidateExclusiveObjectDomain::CirclePackage {
298 reference: reference.clone(),
299 },
300 )),
301 SharedLiveSetObjectDomain::StoredBlob => None,
302 SharedLiveSetObjectDomain::StoreSnapshotImage { .. } => None,
303 SharedLiveSetObjectDomain::StoreMembershipRollup { .. } => None,
304 SharedLiveSetObjectDomain::CircleBootstrapImage { .. } => None,
305 };
306 if let Some((family, domain)) = package_domain {
307 let identity = CandidateExclusiveTarget {
308 family,
309 domain,
310 semantic_hash: record.identity.semantic_hash,
311 object: record.identity.object.clone(),
312 };
313 let payloads = record.payloads.clone();
314 *self = Self::CandidateExclusive(CandidateObjectRecord {
315 identity,
316 payloads,
317 state: CandidateObjectState::CleanupPending { former_candidates },
318 });
319 } else {
320 record.state = OwnedObjectState::RetirementPending { former_candidates };
321 }
322 Ok(())
323 }
324
325 pub fn merge_snapshot_owner(
326 &mut self,
327 stored: &crate::blob::locator::StoredBlobRef,
328 owner: SnapshotObjectOwner,
329 ) -> Result<(), RemoteObjectRecordError> {
330 let Self::SharedLiveSet(record) = self else {
331 return Err(RemoteObjectRecordError::DomainMismatch);
332 };
333 let locator_bytes = stored.locator().to_bytes();
334 if record.identity.domain != SharedLiveSetObjectDomain::StoredBlob
335 || record.identity.semantic_hash != ObjectHash::digest(&locator_bytes)
336 || record.identity.object != *stored.object()
337 || record.payloads.carried_locator_bytes() != Some(locator_bytes.as_slice())
338 {
339 return Err(RemoteObjectRecordError::StoredReferenceMismatch);
340 }
341 merge_shared_owner(&mut record.state, SharedObjectOwner::Snapshot(owner))?;
342 self.validate()
343 }
344
345 pub fn replace_snapshot_owners_for_image(
348 &mut self,
349 owner: Option<&SnapshotObjectOwner>,
350 pending_store_snapshots: &BTreeSet<crate::objects::ObjectSlot>,
351 ) -> Result<(), RemoteObjectRecordError> {
352 if let Self::SharedLiveSet(record) = self {
353 if let OwnedObjectState::UploadedVerified { ownership } = &mut record.state {
354 ownership.activated.retain(|owner| match owner {
355 SharedObjectOwner::Snapshot(SnapshotObjectOwner::Store { metadata_slot }) => {
356 pending_store_snapshots.contains(metadata_slot)
357 }
358 SharedObjectOwner::Snapshot(SnapshotObjectOwner::Circle { .. }) => false,
359 _ => true,
360 });
361 if let Some(owner) = owner {
362 ownership
363 .activated
364 .insert(SharedObjectOwner::Snapshot(owner.clone()));
365 }
366 }
367 }
368 self.validate()
369 }
370
371 pub fn retire_superseded_store_snapshot_ownership(
375 &mut self,
376 metadata_slot: &crate::objects::ObjectSlot,
377 superseded: &BTreeSet<crate::objects::ObjectSlot>,
378 ) -> Result<(), RemoteObjectRecordError> {
379 let Self::SharedLiveSet(record) = self else {
380 return Err(RemoteObjectRecordError::DomainMismatch);
381 };
382 let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
383 return Err(RemoteObjectRecordError::InvalidActivation);
384 };
385 let current = SharedObjectOwner::Snapshot(SnapshotObjectOwner::Store {
386 metadata_slot: metadata_slot.clone(),
387 });
388 let published_blob = record.identity.domain == SharedLiveSetObjectDomain::StoredBlob
389 && ownership
390 .activated
391 .iter()
392 .any(|owner| matches!(owner, SharedObjectOwner::StoreCommit(_)));
393 if (!ownership.activated.contains(¤t) && !published_blob)
394 || superseded.contains(metadata_slot)
395 {
396 return Err(RemoteObjectRecordError::CandidateOwnerMismatch);
397 }
398 ownership.activated.retain(|owner| match owner {
399 SharedObjectOwner::Snapshot(SnapshotObjectOwner::Store { metadata_slot }) => {
400 !superseded.contains(metadata_slot)
401 }
402 _ => true,
403 });
404 self.validate()
405 }
406}