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(deny_unknown_fields)]
166pub struct SnapshotObjectOwner {
167    pub activation: StreamActivationId,
168    pub generation: u64,
169}
170
171impl RemoteObjectRecord {
172    pub fn merge_blob_activation(
173        &mut self,
174        stored: &crate::blob::locator::StoredBlobRef,
175        owner: &StoreBatchCommitRef,
176    ) -> Result<(), RemoteObjectRecordError> {
177        let Self::SharedLiveSet(record) = self else {
178            return Err(RemoteObjectRecordError::DomainMismatch);
179        };
180        let locator_bytes = stored.locator().to_bytes();
181        if record.identity.domain != SharedLiveSetObjectDomain::StoredBlob
182            || record.identity.semantic_hash != ObjectHash::digest(&locator_bytes)
183            || record.identity.object != *stored.object()
184            || record.payloads.carried_locator_bytes() != Some(locator_bytes.as_slice())
185        {
186            return Err(RemoteObjectRecordError::StoredReferenceMismatch);
187        }
188        merge_store_commit_owner(&mut record.state, owner);
189        self.validate()
190    }
191
192    /// Add one more snapshot generation to a membership rollup's owners.
193    ///
194    /// A rollup is content-addressed over the membership frontier, so a
195    /// generation published while membership stood still names the object an
196    /// earlier one already owns. Both own it; reclaim deletes it when the last
197    /// owner goes.
198    pub fn merge_snapshot_ownership(
199        &mut self,
200        rollup: &crate::store_commit::MembershipRollupRef,
201        owner: SnapshotObjectOwner,
202    ) -> Result<(), RemoteObjectRecordError> {
203        let Self::SharedLiveSet(record) = self else {
204            return Err(RemoteObjectRecordError::DomainMismatch);
205        };
206        if !matches!(
207            &record.identity.domain,
208            SharedLiveSetObjectDomain::StoreMembershipRollup { reference } if reference == rollup
209        ) || record.identity.semantic_hash != rollup.rollup_hash
210            || record.identity.object != rollup.object
211        {
212            return Err(RemoteObjectRecordError::StoredReferenceMismatch);
213        }
214        let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
215            return Err(RemoteObjectRecordError::DomainMismatch);
216        };
217        ownership
218            .activated
219            .insert(SharedObjectOwner::Snapshot(owner));
220        self.validate()
221    }
222
223    pub fn merge_package_activation(
224        &mut self,
225        domain: &SharedLiveSetObjectDomain,
226        package: &crate::audience_package::AudiencePackage,
227        owner: &StoreBatchCommitRef,
228    ) -> Result<(), RemoteObjectRecordError> {
229        if !matches!(
230            domain,
231            SharedLiveSetObjectDomain::StorePackage { .. }
232                | SharedLiveSetObjectDomain::CirclePackage { .. }
233        ) {
234            return Err(RemoteObjectRecordError::DomainMismatch);
235        }
236        let Self::SharedLiveSet(record) = self else {
237            return Err(RemoteObjectRecordError::DomainMismatch);
238        };
239        let canonical_semantic_bytes = package.to_bytes();
240        if &record.identity.domain != domain
241            || record.identity.semantic_hash != ObjectHash::digest(&canonical_semantic_bytes)
242            || record.identity.object != *domain.package_object()?
243            || matches!(record.payloads, RemoteObjectPayloads::RowBlob { .. })
244        {
245            return Err(RemoteObjectRecordError::StoredReferenceMismatch);
246        }
247        merge_store_commit_owner(&mut record.state, owner);
248        self.validate()
249    }
250
251    pub fn merge_retained_replay_owner(
252        &mut self,
253        owner: RetainedReplayOwner,
254    ) -> Result<(), RemoteObjectRecordError> {
255        let Self::SharedLiveSet(record) = self else {
256            return Err(RemoteObjectRecordError::DomainMismatch);
257        };
258        let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
259            return Err(RemoteObjectRecordError::InvalidActivation);
260        };
261        ownership
262            .activated
263            .insert(SharedObjectOwner::RetainedReplay(owner));
264        self.validate()
265    }
266
267    pub fn remove_all_retained_replay_owners(&mut self) -> Result<(), RemoteObjectRecordError> {
268        let Self::SharedLiveSet(record) = self else {
269            return Ok(());
270        };
271        let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
272            return Ok(());
273        };
274        ownership
275            .activated
276            .retain(|owner| !matches!(owner, SharedObjectOwner::RetainedReplay(_)));
277        self.validate()
278    }
279
280    pub fn remove_retained_replay_owner(
281        &mut self,
282        owner: &RetainedReplayOwner,
283    ) -> Result<(), RemoteObjectRecordError> {
284        let Self::SharedLiveSet(record) = self else {
285            return Err(RemoteObjectRecordError::DomainMismatch);
286        };
287        let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
288            return Err(RemoteObjectRecordError::InvalidActivation);
289        };
290        if !ownership
291            .activated
292            .remove(&SharedObjectOwner::RetainedReplay(owner.clone()))
293        {
294            return Err(RemoteObjectRecordError::CandidateOwnerMismatch);
295        }
296        self.retire_unowned_shared_live_set()?;
297        self.validate()
298    }
299
300    fn retire_unowned_shared_live_set(&mut self) -> Result<(), RemoteObjectRecordError> {
301        let Self::SharedLiveSet(record) = self else {
302            return Err(RemoteObjectRecordError::DomainMismatch);
303        };
304        let OwnedObjectState::UploadedVerified { ownership } = &record.state else {
305            return Err(RemoteObjectRecordError::InvalidActivation);
306        };
307        if !ownership.pending.is_empty() || !ownership.activated.is_empty() {
308            return Ok(());
309        }
310        if ownership.nonactivated.is_empty() {
311            return Err(RemoteObjectRecordError::EmptyOwnership);
312        }
313        let former_candidates = ownership.nonactivated.clone();
314        let package_domain = match &record.identity.domain {
315            SharedLiveSetObjectDomain::StorePackage { reference } => Some((
316                reference.candidate_family,
317                CandidateExclusiveObjectDomain::StorePackage {
318                    reference: reference.clone(),
319                },
320            )),
321            SharedLiveSetObjectDomain::CirclePackage { reference } => Some((
322                reference.package.candidate_family,
323                CandidateExclusiveObjectDomain::CirclePackage {
324                    reference: reference.clone(),
325                },
326            )),
327            SharedLiveSetObjectDomain::StoredBlob => None,
328            SharedLiveSetObjectDomain::StoreSnapshotImage { .. } => None,
329            SharedLiveSetObjectDomain::StoreMembershipRollup { .. } => None,
330            SharedLiveSetObjectDomain::CircleBootstrapImage { .. } => None,
331        };
332        if let Some((family, domain)) = package_domain {
333            let identity = CandidateExclusiveTarget {
334                family,
335                domain,
336                semantic_hash: record.identity.semantic_hash,
337                object: record.identity.object.clone(),
338            };
339            let payloads = record.payloads.clone();
340            *self = Self::CandidateExclusive(CandidateObjectRecord {
341                identity,
342                payloads,
343                state: CandidateObjectState::CleanupPending { former_candidates },
344            });
345        } else {
346            record.state = OwnedObjectState::RetirementPending { former_candidates };
347        }
348        Ok(())
349    }
350
351    pub fn retract_activated_candidate(
352        &mut self,
353        nonactivation: CandidateNonactivation,
354        head_nonactivation: Option<&VerifiedCandidateHeadNonactivation>,
355    ) -> Result<Option<ProtocolInertObject>, RemoteObjectRecordError> {
356        nonactivation.validate()?;
357        if !matches!(
358            nonactivation.proof,
359            CandidateNonactivationProof::AuthorExclusion { .. }
360                | CandidateNonactivationProof::MergeMembershipGrantRevocation { .. }
361        ) {
362            return Err(RemoteObjectRecordError::InvalidProof(
363                "activated candidate retraction requires terminal author exclusion".to_string(),
364            ));
365        }
366        let candidate = nonactivation.reference()?;
367        match self {
368            Self::RetainedAuthority(record) => {
369                let RetainedAuthorityObjectState::UploadedVerified { ownership } =
370                    &mut record.state
371                else {
372                    return Err(RemoteObjectRecordError::InvalidActivation);
373                };
374                if !ownership.activated.remove(&candidate) {
375                    ensure_candidate_nonactivation(&ownership.nonactivated, &candidate)?;
376                    return Ok(None);
377                }
378                ownership.nonactivated.push(nonactivation.clone());
379                let is_head = matches!(
380                    record.identity.domain,
381                    RetainedAuthorityObjectDomain::DeviceHead { .. }
382                );
383                match (is_head, head_nonactivation) {
384                    (true, Some(head_nonactivation))
385                        if head_nonactivation.candidate == candidate
386                            && matches!(
387                                &head_nonactivation.head,
388                                VerifiedCandidateHead::ExactLateCandidate { object }
389                                    if object == &record.identity.object
390                            ) => {}
391                    (true, _) => {
392                        return Err(RemoteObjectRecordError::InvalidProof(
393                            "uploaded retracted candidate head lacks exact presence evidence"
394                                .to_string(),
395                        ));
396                    }
397                    (false, None) => {}
398                    (false, Some(_)) => {
399                        return Err(RemoteObjectRecordError::InvalidProof(
400                            "candidate-head evidence reached a non-head activated object"
401                                .to_string(),
402                        ));
403                    }
404                }
405                if !ownership.pending.is_empty() || !ownership.activated.is_empty() {
406                    self.validate()?;
407                    return Ok(None);
408                }
409                if matches!(
410                    &record.identity.domain,
411                    RetainedAuthorityObjectDomain::Commit { reference }
412                        if reference == &candidate
413                ) {
414                    let payloads = record.payloads.clone();
415                    *self = Self::CandidateCommit(CandidateCommitRecord {
416                        identity: candidate,
417                        semantic_hash: record.identity.semantic_hash,
418                        payloads,
419                        state: CandidateCommitState::CleanupPending {
420                            proof: nonactivation.proof,
421                        },
422                    });
423                    self.validate()?;
424                    return Ok(None);
425                }
426                ProtocolInertObject::new(record.identity.clone(), ownership.nonactivated.clone())
427                    .map(Some)
428            }
429            Self::SharedLiveSet(record) => {
430                if head_nonactivation.is_some() {
431                    return Err(RemoteObjectRecordError::InvalidProof(
432                        "candidate-head evidence reached a shared activated object".to_string(),
433                    ));
434                }
435                let OwnedObjectState::UploadedVerified { ownership } = &mut record.state else {
436                    return Err(RemoteObjectRecordError::InvalidActivation);
437                };
438                if !ownership
439                    .activated
440                    .remove(&SharedObjectOwner::StoreCommit(candidate.clone()))
441                {
442                    ensure_candidate_nonactivation(&ownership.nonactivated, &candidate)?;
443                    return Ok(None);
444                }
445                ownership.nonactivated.push(nonactivation);
446                self.retire_unowned_shared_live_set()?;
447                self.validate()?;
448                Ok(None)
449            }
450            Self::CandidateCommit(_) | Self::CandidateExclusive(_) => {
451                Err(RemoteObjectRecordError::InvalidActivation)
452            }
453        }
454    }
455
456    pub fn merge_snapshot_owner(
457        &mut self,
458        stored: &crate::blob::locator::StoredBlobRef,
459        owner: SnapshotObjectOwner,
460    ) -> Result<(), RemoteObjectRecordError> {
461        let Self::SharedLiveSet(record) = self else {
462            return Err(RemoteObjectRecordError::DomainMismatch);
463        };
464        let locator_bytes = stored.locator().to_bytes();
465        if record.identity.domain != SharedLiveSetObjectDomain::StoredBlob
466            || record.identity.semantic_hash != ObjectHash::digest(&locator_bytes)
467            || record.identity.object != *stored.object()
468            || record.payloads.carried_locator_bytes() != Some(locator_bytes.as_slice())
469        {
470            return Err(RemoteObjectRecordError::StoredReferenceMismatch);
471        }
472        merge_shared_owner(&mut record.state, SharedObjectOwner::Snapshot(owner))?;
473        self.validate()
474    }
475}