Skip to main content

coven_protocol/remote_object/
ownership.rs

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    /// Reassign inherited snapshot ownership in an exported database copy.
346    /// Live records keep their owners through the merge methods instead.
347    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    /// Retire superseded Store snapshot leases at an accepted successor. A
372    /// reused object remains owned by that successor; a dropped blob retains
373    /// its exact original publication provenance until physical reclaim.
374    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(&current) && !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}