Skip to main content

coven_protocol/provider/
admin.rs

1use super::probe::*;
2use super::*;
3
4#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
5#[serde(deny_unknown_fields)]
6pub struct ProviderCapabilityProof {
7    pub exact_slots: ExactSlotProbeReceipt,
8}
9
10#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
11#[serde(deny_unknown_fields)]
12pub struct FounderProviderAdminGrant {
13    pub grant_id: ProviderAdminGrantId,
14    pub provider: ProviderDeviceBinding,
15    pub access: ProviderAccessLocator,
16    pub capability: ProviderCapabilityProof,
17}
18
19impl FounderProviderAdminGrant {
20    #[cfg(any(test, feature = "test-utils"))]
21    pub fn from_test_label(label: &str) -> Self {
22        let probe_id =
23            ProviderProbeId::from_bytes(*ObjectHash::digest(label.as_bytes()).as_bytes());
24        let slot = ObjectSlot::logical(format!("store-v1/test/{label}/provider-probe/exact"))
25            .expect("valid exact-probe test slot");
26        let first = probe_payload(&probe_id, ProbePayloadLabel::ExactCreateFirst);
27        let second = probe_payload(&probe_id, ProbePayloadLabel::ExactCreateSecond);
28        let accepted =
29            ExactObjectRef::new(slot.clone(), first.len() as u64, ObjectHash::digest(&first));
30        let lost_slot = ObjectSlot::logical(format!(
31            "store-v1/test/{label}/provider-probe/lost-response"
32        ))
33        .expect("valid lost-response test slot");
34        let lost_payload = probe_payload(&probe_id, ProbePayloadLabel::LostResponse);
35        let lost_ref = ExactObjectRef::new(
36            lost_slot.clone(),
37            lost_payload.len() as u64,
38            ObjectHash::digest(&lost_payload),
39        );
40        let conditional_slot =
41            ObjectSlot::logical(format!("store-v1/test/{label}/provider-probe/conditional"))
42                .expect("valid conditional-update test slot");
43        let device = ProviderDeviceBinding {
44            principal: crate::objects::ProviderPrincipalId::CustomS3Credential {
45                access_key_id_hash: ObjectHash::digest(format!("{label} access key").as_bytes()),
46            },
47        };
48        let store = StoreProviderBinding::S3 {
49            endpoint: crate::objects::S3EndpointBinding::Custom {
50                origin: "https://test.invalid".to_string(),
51            },
52            region: "test-region".to_string(),
53            bucket: format!("{label}-bucket"),
54            key_prefix: None,
55        };
56        let transcript = ExactSlotProbeTranscript {
57            probe_id,
58            logical_key: slot.logical_key().to_string(),
59            slot,
60            contenders: [
61                ProbeCreateAttempt {
62                    payload_hash: ObjectHash::digest(&first),
63                    outcome: ProbeCreateOutcome::Created,
64                },
65                ProbeCreateAttempt {
66                    payload_hash: ObjectHash::digest(&second),
67                    outcome: ProbeCreateOutcome::RejectedOccupied,
68                },
69            ],
70            accepted: accepted.clone(),
71            full_read_hash: accepted.stored_hash(),
72            range: ProbeRangeReceipt {
73                start: PROBE_RANGE_START,
74                end: PROBE_RANGE_END,
75                bytes_hash: ObjectHash::digest(
76                    &first[PROBE_RANGE_START as usize..PROBE_RANGE_END as usize],
77                ),
78            },
79            conditional: crate::provider::test_fixtures::test_conditional_receipt(
80                probe_id,
81                conditional_slot,
82            ),
83            lost_response: LostResponseProbeReceipt {
84                logical_key: lost_slot.logical_key().to_string(),
85                slot: lost_slot,
86                payload_hash: ObjectHash::digest(&lost_payload),
87                settled: lost_ref,
88                readback_hash: ObjectHash::digest(&lost_payload),
89            },
90        };
91        Self {
92            grant_id: ProviderAdminGrantId(ObjectHash::digest(
93                format!("{label} provider admin grant").as_bytes(),
94            )),
95            provider: device.clone(),
96            access: ProviderAccessLocator::S3SharedCredentialGeneration {
97                generation: 1,
98                access_key_id_hash: ObjectHash::digest(format!("{label} access key").as_bytes()),
99            },
100            capability: ProviderCapabilityProof {
101                exact_slots: ExactSlotProbeReceipt::from_transcript(transcript, &store, &device),
102            },
103        }
104    }
105}
106
107impl ProviderCapabilityProof {
108    pub fn verify(
109        &self,
110        store: &StoreProviderBinding,
111        device: &ProviderDeviceBinding,
112    ) -> Result<(), ProviderProbeError> {
113        self.exact_slots.verify(store, device)
114    }
115}
116
117#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
118#[serde(deny_unknown_fields)]
119pub struct ProviderAdminGrantRecord {
120    pub grant_id: ProviderAdminGrantId,
121    pub administrator: StoreDeviceRegistrationRef,
122    pub provider: ProviderDeviceBinding,
123    pub access: ProviderAccessLocator,
124    pub capability: ProviderCapabilityProof,
125    pub created_at: ProviderAdminGrantOrigin,
126}
127
128#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
129#[serde(rename_all = "snake_case", deny_unknown_fields)]
130pub enum ProviderAdminGrantOrigin {
131    Founder { root: StoreRootRef },
132    Membership { coord: MembershipCoord },
133}
134
135#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
136#[serde(deny_unknown_fields)]
137pub struct ProviderAdminMembershipChange {
138    pub change: ProviderAdminChange,
139    #[serde(with = "ordered_owner_barriers")]
140    pub owner_barriers: BTreeMap<MembershipGrantId, OwnerStreamBarrier>,
141}
142
143#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
144#[serde(rename_all = "snake_case", deny_unknown_fields)]
145pub enum ProviderAdminChange {
146    Set {
147        administrator: StoreDeviceRegistrationRef,
148        provider: ProviderDeviceBinding,
149        access: ProviderAccessLocator,
150        capability: ProviderCapabilityProof,
151        grant_id: ProviderAdminGrantId,
152        replaces: BTreeSet<ProviderAdminGrantId>,
153    },
154    Remove {
155        removes: BTreeSet<ProviderAdminGrantId>,
156    },
157}
158
159#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
160#[serde(deny_unknown_fields)]
161pub struct ProviderAdminState {
162    records: BTreeMap<ProviderAdminGrantId, ProviderAdminGrantRecord>,
163    tombstones: BTreeSet<ProviderAdminGrantId>,
164}
165
166#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
167#[serde(deny_unknown_fields)]
168pub struct ProviderAdminBranch {
169    pub heads: Vec<MembershipCoord>,
170    pub state: ProviderAdminState,
171}
172
173#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
174#[serde(deny_unknown_fields)]
175pub struct ProviderAdminConflict {
176    pub raw_heads: Vec<MembershipCoord>,
177    pub cyclic_sources: Vec<MembershipCoord>,
178    pub involved_grants: BTreeSet<ProviderAdminGrantId>,
179    pub maximal_valid_branches: Vec<ProviderAdminBranch>,
180    pub combined: ProviderAdminState,
181}
182
183#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
184#[serde(rename_all = "snake_case", deny_unknown_fields)]
185pub enum ProviderAdminResolution {
186    Resolved(ProviderAdminState),
187    RevocationConflict(ProviderAdminConflict),
188}
189
190impl ProviderAdminResolution {
191    pub fn combined_state(&self) -> &ProviderAdminState {
192        match self {
193            Self::Resolved(state) => state,
194            Self::RevocationConflict(conflict) => &conflict.combined,
195        }
196    }
197
198    pub fn state_hash(&self) -> ObjectHash {
199        ObjectHash::digest(&domain_json(b"coven.provider-admin-resolution.v1\0", self))
200    }
201}
202
203impl ProviderAdminState {
204    pub fn founder(grant: ProviderAdminGrantRecord) -> Self {
205        let grant_id = grant.grant_id.clone();
206        Self {
207            records: BTreeMap::from([(grant_id.clone(), grant)]),
208            tombstones: BTreeSet::new(),
209        }
210    }
211
212    pub fn founder_from_root(
213        root: StoreRootRef,
214        administrator: StoreDeviceRegistrationRef,
215        grant: &FounderProviderAdminGrant,
216    ) -> Self {
217        Self::founder(ProviderAdminGrantRecord {
218            grant_id: grant.grant_id.clone(),
219            administrator,
220            provider: grant.provider.clone(),
221            access: grant.access.clone(),
222            capability: grant.capability.clone(),
223            created_at: ProviderAdminGrantOrigin::Founder { root },
224        })
225    }
226
227    pub fn authorizes(
228        &self,
229        grant_id: &ProviderAdminGrantId,
230        administrator: &StoreDeviceRegistrationRef,
231    ) -> bool {
232        !self.tombstones.contains(grant_id)
233            && self
234                .records
235                .get(grant_id)
236                .is_some_and(|record| &record.administrator == administrator)
237    }
238
239    pub fn records(&self) -> &BTreeMap<ProviderAdminGrantId, ProviderAdminGrantRecord> {
240        &self.records
241    }
242
243    pub fn active(&self) -> BTreeSet<ProviderAdminGrantId> {
244        self.records
245            .keys()
246            .filter(|grant_id| !self.tombstones.contains(*grant_id))
247            .cloned()
248            .collect()
249    }
250
251    pub fn tombstones(&self) -> &BTreeSet<ProviderAdminGrantId> {
252        &self.tombstones
253    }
254
255    pub fn apply(
256        &mut self,
257        change: ProviderAdminChange,
258        origin: ProviderAdminGrantOrigin,
259    ) -> Result<(), ProviderAdminReducerError> {
260        let mut next = self.clone();
261        next.apply_unchecked(change, origin)?;
262        if next.active().is_empty() {
263            return Err(ProviderAdminReducerError::NoEffectiveAdministrator);
264        }
265        *self = next;
266        Ok(())
267    }
268
269    fn apply_unchecked(
270        &mut self,
271        change: ProviderAdminChange,
272        origin: ProviderAdminGrantOrigin,
273    ) -> Result<(), ProviderAdminReducerError> {
274        match change {
275            ProviderAdminChange::Set {
276                administrator,
277                provider,
278                access,
279                capability,
280                grant_id,
281                replaces,
282            } => {
283                let record = ProviderAdminGrantRecord {
284                    grant_id: grant_id.clone(),
285                    administrator,
286                    provider,
287                    access,
288                    capability,
289                    created_at: origin,
290                };
291                if let Some(existing) = self.records.get(&grant_id) {
292                    if existing != &record {
293                        return Err(ProviderAdminReducerError::GrantIdReuse);
294                    }
295                    if !replaces.iter().all(|id| self.tombstones.contains(id)) {
296                        return Err(ProviderAdminReducerError::UnknownReplacement);
297                    }
298                    return Ok(());
299                }
300                if !replaces
301                    .iter()
302                    .all(|id| self.records.contains_key(id) && !self.tombstones.contains(id))
303                {
304                    return Err(ProviderAdminReducerError::UnknownReplacement);
305                }
306                for replaced in replaces {
307                    self.tombstones.insert(replaced);
308                }
309                self.records.insert(grant_id, record);
310            }
311            ProviderAdminChange::Remove { removes } => {
312                if removes.is_empty()
313                    || !removes
314                        .iter()
315                        .all(|id| self.records.contains_key(id) || self.tombstones.contains(id))
316                {
317                    return Err(ProviderAdminReducerError::UnknownRemoval);
318                }
319                for removed in removes {
320                    self.tombstones.insert(removed);
321                }
322            }
323        }
324        Ok(())
325    }
326
327    pub(crate) fn apply_membership_change(
328        &mut self,
329        change: ProviderAdminMembershipChange,
330        origin: ProviderAdminGrantOrigin,
331    ) -> Result<(), ProviderAdminReducerError> {
332        if !matches!(origin, ProviderAdminGrantOrigin::Membership { .. }) {
333            return Err(ProviderAdminReducerError::PolicyOriginMismatch);
334        }
335        self.apply(change.change, origin)
336    }
337
338    pub fn state_hash(&self) -> ObjectHash {
339        ObjectHash::digest(&domain_json(
340            b"coven.provider-admin-state.v1\0",
341            &(self.records(), self.tombstones()),
342        ))
343    }
344
345    pub fn merge(
346        states: impl IntoIterator<Item = Self>,
347    ) -> Result<Self, ProviderAdminReducerError> {
348        let mut records = BTreeMap::new();
349        let mut tombstones = BTreeSet::new();
350        for state in states {
351            for (grant_id, record) in state.records {
352                if records
353                    .insert(grant_id.clone(), record.clone())
354                    .is_some_and(|current| current != record)
355                {
356                    return Err(ProviderAdminReducerError::GrantIdReuse);
357                }
358            }
359            tombstones.extend(state.tombstones);
360        }
361        Ok(Self {
362            records,
363            tombstones,
364        })
365    }
366
367    pub(crate) fn reduce_merge(
368        genesis: &Self,
369        entries: &[MembershipEntry],
370        included: &BTreeSet<MembershipCoord>,
371    ) -> Result<ProviderAdminResolution, ProviderAdminReducerError> {
372        let by_coord = entries
373            .iter()
374            .filter(|entry| included.contains(&entry.coord()))
375            .map(|entry| (entry.coord(), entry))
376            .collect::<BTreeMap<_, _>>();
377        let mut states = BTreeMap::<MembershipCoord, Self>::new();
378        let mut pending = by_coord.keys().cloned().collect::<BTreeSet<_>>();
379        while !pending.is_empty() {
380            let ready = pending.iter().find(|coord| {
381                let entry = by_coord[*coord];
382                let predecessor = (entry.seq > 1)
383                    .then(|| stream_predecessor(&by_coord, entry))
384                    .flatten();
385                (entry.seq == 1 || predecessor.is_some_and(|value| states.contains_key(value)))
386                    && entry
387                        .dependencies
388                        .iter()
389                        .filter(|dependency| included.contains(*dependency))
390                        .all(|dependency| states.contains_key(dependency))
391            });
392            let Some(coord) = ready.cloned() else {
393                if pending.iter().any(|coord| {
394                    let entry = by_coord[coord];
395                    entry.seq > 1 && stream_predecessor(&by_coord, entry).is_none()
396                }) {
397                    return Err(ProviderAdminReducerError::MissingPredecessor);
398                }
399                return Err(ProviderAdminReducerError::CausalCycle);
400            };
401            let entry = by_coord[&coord];
402            let mut causal_states = entry
403                .dependencies
404                .iter()
405                .filter_map(|dependency| states.get(dependency).cloned())
406                .collect::<Vec<_>>();
407            if entry.seq > 1 {
408                if let Some(predecessor) = stream_predecessor(&by_coord, entry) {
409                    if !entry.dependencies.contains(predecessor) {
410                        causal_states.push(states[predecessor].clone());
411                    }
412                }
413            }
414            let mut state = if causal_states.is_empty() {
415                genesis.clone()
416            } else {
417                Self::merge(causal_states)?
418            };
419            if let Some(change) = entry.provider_admin.clone() {
420                state.apply_membership_change(
421                    change,
422                    ProviderAdminGrantOrigin::Membership {
423                        coord: coord.clone(),
424                    },
425                )?;
426            }
427            states.insert(coord.clone(), state);
428            pending.remove(&coord);
429        }
430        let raw_heads = by_coord
431            .keys()
432            .filter(|coord| {
433                !by_coord.values().any(|entry| {
434                    entry.dependencies.contains(*coord)
435                        || (entry.seq == coord.seq + 1
436                            && entry.author_pubkey == coord.author_pubkey
437                            && entry.author_owner_grant == coord.author_owner_grant
438                            && entry.stream_id == coord.stream_id
439                            && entry.previous_hash == Some(coord.entry_hash))
440                })
441            })
442            .cloned()
443            .collect::<Vec<_>>();
444        let combined =
445            Self::merge(std::iter::once(genesis.clone()).chain(states.values().cloned()))?;
446        if !combined.active().is_empty() {
447            return Ok(ProviderAdminResolution::Resolved(combined));
448        }
449        if raw_heads.len() > 12 {
450            return Err(ProviderAdminReducerError::ConflictTooWide(raw_heads.len()));
451        }
452        let head_states = raw_heads
453            .iter()
454            .map(|head| (head.clone(), states[head].clone()))
455            .collect::<Vec<_>>();
456        let mut valid = Vec::<ProviderAdminBranch>::new();
457        for mask in 1usize..(1usize << head_states.len()) {
458            let heads = head_states
459                .iter()
460                .enumerate()
461                .filter(|(index, _)| mask & (1usize << index) != 0)
462                .map(|(_, (head, _))| head.clone())
463                .collect::<Vec<_>>();
464            let state = Self::merge(
465                head_states
466                    .iter()
467                    .enumerate()
468                    .filter(|(index, _)| mask & (1usize << index) != 0)
469                    .map(|(_, (_, state))| state.clone()),
470            )?;
471            if !state.active().is_empty() {
472                valid.push(ProviderAdminBranch { heads, state });
473            }
474        }
475        let valid_head_sets = valid
476            .iter()
477            .map(|branch| branch.heads.iter().cloned().collect::<BTreeSet<_>>())
478            .collect::<Vec<_>>();
479        let maximal_valid_branches = valid
480            .into_iter()
481            .enumerate()
482            .filter(|(index, _)| {
483                !valid_head_sets.iter().enumerate().any(|(other, heads)| {
484                    other != *index && valid_head_sets[*index].is_subset(heads)
485                })
486            })
487            .map(|(_, branch)| branch)
488            .collect();
489        let mut cyclic_sources = Vec::new();
490        let mut involved_grants = BTreeSet::new();
491        for (coord, entry) in &by_coord {
492            if let Some(ProviderAdminMembershipChange {
493                change: ProviderAdminChange::Remove { removes },
494                ..
495            }) = &entry.provider_admin
496            {
497                cyclic_sources.push(coord.clone());
498                involved_grants.extend(removes.iter().cloned());
499            }
500        }
501        cyclic_sources.sort();
502        Ok(ProviderAdminResolution::RevocationConflict(
503            ProviderAdminConflict {
504                raw_heads,
505                cyclic_sources,
506                involved_grants,
507                maximal_valid_branches,
508                combined,
509            },
510        ))
511    }
512}
513
514#[derive(Debug, thiserror::Error, PartialEq, Eq)]
515pub enum ProviderAdminReducerError {
516    #[error("provider administrator grant id was reused with different facts")]
517    GrantIdReuse,
518    #[error("provider administrator replacement names an inactive grant")]
519    UnknownReplacement,
520    #[error("provider administrator removal names an inactive grant")]
521    UnknownRemoval,
522    #[error("provider administrator change leaves no effective administrator")]
523    NoEffectiveAdministrator,
524    #[error("provider administrator change policy does not match its derived origin")]
525    PolicyOriginMismatch,
526    #[error("provider administrator causal history is missing an exact stream predecessor")]
527    MissingPredecessor,
528    #[error("provider administrator causal history contains a cycle")]
529    CausalCycle,
530    #[error("provider administrator revocation conflict has {0} heads, exceeding 12")]
531    ConflictTooWide(usize),
532}
533
534/// The in-stream predecessor of `entry` among the included coordinates: the
535/// same author stream at the previous sequence, carrying the hash the entry
536/// links back to.
537fn stream_predecessor<'coords>(
538    by_coord: &'coords BTreeMap<MembershipCoord, &MembershipEntry>,
539    entry: &MembershipEntry,
540) -> Option<&'coords MembershipCoord> {
541    by_coord.keys().find(|candidate| {
542        candidate.author_pubkey == entry.author_pubkey
543            && candidate.author_owner_grant == entry.author_owner_grant
544            && candidate.stream_id == entry.stream_id
545            && candidate.seq + 1 == entry.seq
546            && Some(candidate.entry_hash) == entry.previous_hash
547    })
548}