Skip to main content

coven_database/store/store_session/
membership_mutations.rs

1use super::*;
2use crate::store::StoreSession;
3use crate::*;
4use coven_protocol::store_commit::ObjectHash;
5use rusqlite::OptionalExtension;
6use std::collections::BTreeSet;
7
8impl StoreSession<'_> {
9    fn outbound_membership_mutation(
10        &mut self,
11    ) -> Result<Option<DurableMembershipMutation>, DbError> {
12        let conn = self.conn;
13        conn.query_row(
14            "SELECT intent_hash, plan_bytes, progress_bytes \
15             FROM outbound_membership_mutation WHERE singleton = 1",
16            [],
17            |row| {
18                Ok((
19                    row.get::<_, String>(0)?,
20                    row.get::<_, Vec<u8>>(1)?,
21                    row.get::<_, Vec<u8>>(2)?,
22                ))
23            },
24        )
25        .optional()
26        .map_err(DbError::from)?
27        .map(|(hash, plan_bytes, progress_bytes)| {
28            let intent_hash: ObjectHash = hash
29                .parse()
30                .map_err(|error| DbError::context("membership intent hash", error))?;
31            if ObjectHash::digest(&plan_bytes) != intent_hash {
32                return Err(DbError::Message(
33                    "membership intent hash differs from its exact plan bytes".to_string(),
34                ));
35            }
36            Ok(DurableMembershipMutation {
37                intent_hash,
38                plan_bytes,
39                progress_bytes,
40            })
41        })
42        .transpose()
43    }
44
45    fn select_causal_author_stream(
46        &mut self,
47        key: &str,
48        reusable: &std::collections::BTreeSet<coven_protocol::membership::AuthorStreamId>,
49        candidate: coven_protocol::membership::AuthorStreamId,
50    ) -> Result<coven_protocol::membership::AuthorStreamId, DbError> {
51        let conn = self.conn;
52        let existing = crate::get_protocol_state_on(conn, key)?
53            .map(|value| value.parse().map_err(DbError::from))
54            .transpose()?;
55        if let Some(existing) = existing {
56            if reusable.contains(&existing) {
57                return Ok(existing);
58            }
59        }
60        let selected = reusable.iter().next_back().copied().unwrap_or(candidate);
61        crate::set_protocol_state_on(conn, key, &selected.to_string())?;
62        Ok(selected)
63    }
64
65    fn stage_membership_candidate_mutation(
66        &mut self,
67        plan_bytes: Vec<u8>,
68        progress_bytes: Vec<u8>,
69        remote_objects: Vec<coven_protocol::remote_object::ClosedRemoteObject>,
70        candidate: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
71    ) -> Result<ObjectHash, DbError> {
72        let publication = validate_membership_candidate_objects(&candidate, &remote_objects)?;
73        let pending_rotation_generation = membership_rotation_generation(&publication.entry)?;
74        let active_publication = ActiveStorePublication::for_commit(
75            ActiveStorePublicationOwner::MembershipMutation,
76            &candidate,
77        )?;
78        let conn = self.conn;
79        let intent_hash = ObjectHash::digest(&plan_bytes);
80        let tx = conn.unchecked_transaction().map_err(DbError::from)?;
81        let existing = tx
82            .query_row(
83                "SELECT intent_hash, plan_bytes FROM outbound_membership_mutation \
84                 WHERE singleton = 1",
85                [],
86                |row| Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?)),
87            )
88            .optional()
89            .map_err(DbError::from)?;
90        if let Some((existing_hash, existing_plan)) = existing {
91            if existing_hash != intent_hash.to_string() || existing_plan != plan_bytes {
92                return Err(DbError::Message(
93                    "a different membership mutation is already pending".to_string(),
94                ));
95            }
96            for remote in &remote_objects {
97                let stored = load_remote_object_on(&tx, remote.object_id())?;
98                if stored != **remote {
99                    return Err(DbError::Message(
100                        "persisted membership ownership differs from its durable plan".to_string(),
101                    ));
102                }
103            }
104            if !super::active_store_publication::load_active_store_publication_on(&tx)?
105                .is_some_and(|existing| existing.same_commit_reservation(&active_publication))
106            {
107                return Err(DbError::Message(
108                    "membership candidate differs from its active Store publication".to_string(),
109                ));
110            }
111            super::membership_rotation::stage_pending_rotation_on(
112                &tx,
113                pending_rotation_generation,
114                intent_hash,
115            )?;
116            tx.commit().map_err(DbError::from)?;
117            return Ok(intent_hash);
118        }
119        match super::active_store_publication::claim_active_store_publication_on(
120            &tx,
121            &active_publication,
122        )? {
123            super::active_store_publication::ActiveStorePublicationClaim::Acquired => {}
124            super::active_store_publication::ActiveStorePublicationClaim::AlreadyOwned => {
125                return Err(DbError::Message(
126                    "membership mutation owns publication before its journal".to_string(),
127                ));
128            }
129            super::active_store_publication::ActiveStorePublicationClaim::Occupied(owner) => {
130                return Err(DbError::Message(format!(
131                    "another local Store operation owns publication: {owner:?}"
132                )));
133            }
134        }
135        for remote in &remote_objects {
136            persist_exact_remote_object_on(
137                &tx,
138                self.store_dir,
139                remote,
140                "membership candidate object",
141            )?;
142        }
143        tx.execute(
144            "INSERT INTO outbound_membership_mutation \
145             (singleton, intent_hash, plan_bytes, progress_bytes) \
146             VALUES (1, ?1, ?2, ?3)",
147            rusqlite::params![intent_hash.to_string(), plan_bytes, progress_bytes],
148        )
149        .map_err(DbError::from)?;
150        super::membership_rotation::stage_pending_rotation_on(
151            &tx,
152            pending_rotation_generation,
153            intent_hash,
154        )?;
155        tx.commit().map_err(DbError::from)?;
156        Ok(intent_hash)
157    }
158
159    fn update_membership_mutation_progress(
160        &mut self,
161        intent_hash: ObjectHash,
162        progress_bytes: Vec<u8>,
163    ) -> Result<(), DbError> {
164        let conn = self.conn;
165        let updated = conn
166            .execute(
167                "UPDATE outbound_membership_mutation SET progress_bytes = ?1 \
168                 WHERE singleton = 1 AND intent_hash = ?2",
169                rusqlite::params![progress_bytes, intent_hash.to_string()],
170            )
171            .map_err(DbError::from)?;
172        if updated != 1 {
173            return Err(DbError::Message(
174                "membership mutation ownership row is absent or changed".to_string(),
175            ));
176        }
177        Ok(())
178    }
179
180    fn stage_membership_candidate_abandonment(
181        &mut self,
182        intent_hash: ObjectHash,
183        expected: ActiveStorePublication,
184        original: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
185        abandonment: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
186    ) -> Result<ActiveStorePublication, DbError> {
187        original.validate_closed_shape()?;
188        abandonment.validate_closed_shape()?;
189        let original_target = coven_protocol::store_commit::StoreBatchCommitDeletionTarget {
190            coord: original.reference.coord.clone(),
191            object: original.reference.object.clone(),
192            canonical_signed_bytes: original.commit.to_bytes(),
193        };
194        if expected.owner() != &ActiveStorePublicationOwner::MembershipMutation
195            || expected.commit_reservation()
196                != Some((
197                    &original.commit.write_id,
198                    &original.commit.author_registration,
199                    &original.reference.coord,
200                ))
201            || abandonment.commit.write_id != original.commit.write_id
202            || abandonment.commit.author_registration != original.commit.author_registration
203            || abandonment.reference.coord != original.reference.coord
204            || abandonment.commit.order.predecessor() != original.commit.order.predecessor()
205            || abandonment.commit.abandoned_candidates()
206                != [coven_protocol::store_commit::CandidateCleanupManifest {
207                    candidate: original_target,
208                }]
209            || expected.attempt()?.entry.payload
210                != coven_protocol::store_commit::StorePublicationPayload::Commit(
211                    original.reference.clone(),
212                )
213            || !expected.retired_candidates().is_empty()
214        {
215            return Err(DbError::Message(
216                "membership abandonment differs from its exact reserved candidate".into(),
217            ));
218        }
219        original.prepared_membership_publication()?;
220        let mut replacement = expected.begin_membership_abandonment(abandonment.clone())?;
221        replacement.retain_superseded_entry(expected.attempt()?.reference()?)?;
222        let bytes = abandonment.commit.to_bytes();
223        let remote =
224            RemoteObjectRecord::candidate_commit(abandonment.reference.clone(), &bytes, &bytes)?;
225        let tx = self.conn.unchecked_transaction()?;
226        require_membership_mutation_on(&tx, intent_hash)?;
227        let installed = super::observed_store_publication::load_store_current_publication_on(&tx)?;
228        if installed.record() != &abandonment.publication.previous
229            || installed.observed_version() != Some(&abandonment.publication.previous_version)
230        {
231            return Err(DbError::Message(
232                "membership abandonment does not extend the installed boundary".into(),
233            ));
234        }
235        let entries = super::observed_store_publication::load_store_publication_entries_on(&tx)?;
236        if entries.iter().any(|entry| {
237            matches!(&entry.value.payload,
238                coven_protocol::store_commit::StorePublicationPayload::Commit(candidate)
239                    if candidate.coord == original.reference.coord)
240        }) {
241            return Err(DbError::Message(
242                "accepted membership author position cannot be abandoned".into(),
243            ));
244        }
245        if super::materialized_commit_index::latest_position_for_device_on(
246            &tx,
247            &original.reference.coord.stream_id.to_string(),
248        )?
249        .is_some_and(|tip| tip.coord.sequence >= original.reference.coord.sequence)
250        {
251            return Err(DbError::Message(
252                "membership abandonment cannot consume a covered author position".into(),
253            ));
254        }
255        let previous_attempt = expected.attempt()?.reference()?;
256        let competing = entries.iter().any(|entry| {
257            entry.value.position == previous_attempt.position
258                && entry.prepared.reference() != &previous_attempt.object
259        });
260        let retired = installed
261            .record()
262            .latest_snapshot()
263            .is_some_and(|snapshot| snapshot.publication.position > previous_attempt.position);
264        if !competing && !retired {
265            return Err(DbError::Message(
266                "membership abandonment has no accepted supersession of its previous attempt"
267                    .into(),
268            ));
269        }
270        crate::remote_object_records::validate_remote_object_on(
271            &tx,
272            remote_object_id(&original.reference.object),
273            &original.reference.object,
274            &original.commit.to_bytes(),
275        )?;
276        persist_exact_remote_object_on(&tx, self.store_dir, &remote, "membership abandonment")?;
277        super::active_store_publication::update_active_store_publication_on(
278            &tx,
279            &expected,
280            &replacement,
281        )?;
282        tx.commit()?;
283        Ok(replacement)
284    }
285
286    fn replace_membership_candidate_mutation(
287        &mut self,
288        intent_hash: ObjectHash,
289        expected: ActiveStorePublication,
290        candidate: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
291        plan_bytes: Vec<u8>,
292        progress_bytes: Vec<u8>,
293        remote_objects: Vec<coven_protocol::remote_object::ClosedRemoteObject>,
294    ) -> Result<ObjectHash, DbError> {
295        if expected.owner() != &ActiveStorePublicationOwner::MembershipMutation
296            || !expected.is_awaiting_preparation()
297            || expected.commit_reservation()
298                != Some((
299                    &candidate.commit.write_id,
300                    &candidate.commit.author_registration,
301                    &candidate.reference.coord,
302                ))
303        {
304            return Err(DbError::Message(
305                "replacement membership candidate lacks its exact completed reservation".into(),
306            ));
307        }
308        let publication = validate_membership_candidate_objects(&candidate, &remote_objects)?;
309        let replacement_hash = ObjectHash::digest(&plan_bytes);
310        let tx = self.conn.unchecked_transaction()?;
311        require_membership_mutation_on(&tx, intent_hash)?;
312        let retired = require_retired_membership_candidate_on(&tx, &expected)?;
313        let RetiredStoreCandidateInputs::Membership(original) = &retired.inputs else {
314            unreachable!("retired membership candidate is validated")
315        };
316        use coven_protocol::membership::StoreAuthorityChange;
317        let same_request = match (&original.entry.change, &publication.entry.change) {
318            (
319                StoreAuthorityChange::RemoveMember {
320                    user_pubkey: before,
321                    ..
322                },
323                StoreAuthorityChange::RemoveMember {
324                    user_pubkey: after, ..
325                },
326            ) => before == after,
327            (
328                StoreAuthorityChange::SetMember {
329                    user_pubkey: before,
330                    provider_account_email: old_email,
331                    role: old_role,
332                    ..
333                },
334                StoreAuthorityChange::SetMember {
335                    user_pubkey: after,
336                    provider_account_email: new_email,
337                    role: new_role,
338                    ..
339                },
340            ) => before == after && old_email == new_email && old_role == new_role,
341            _ => false,
342        };
343        if !same_request {
344            return Err(DbError::Message(
345                "replacement changes the retained membership request".into(),
346            ));
347        }
348        let previous_rotation_generation = membership_rotation_generation(&original.entry)?;
349        let pending_rotation_generation = membership_rotation_generation(&publication.entry)?;
350        let replacement = consume_retired_membership_candidate_on(&tx, &expected)?
351            .replace_attempt(candidate.publication.clone())?;
352        let installed = super::observed_store_publication::load_store_current_publication_on(&tx)?;
353        if installed.record() != &candidate.publication.previous
354            || installed.observed_version() != Some(&candidate.publication.previous_version)
355        {
356            return Err(DbError::Message(
357                "replacement membership candidate does not extend the installed boundary".into(),
358            ));
359        }
360        for remote in &remote_objects {
361            persist_exact_remote_object_on(
362                &tx,
363                self.store_dir,
364                remote,
365                "replacement membership candidate object",
366            )?;
367        }
368        if tx.execute(
369            "UPDATE outbound_membership_mutation SET intent_hash = ?1, plan_bytes = ?2, progress_bytes = ?3 \
370             WHERE singleton = 1 AND intent_hash = ?4",
371            rusqlite::params![
372                replacement_hash.to_string(),
373                plan_bytes,
374                progress_bytes,
375                intent_hash.to_string()
376            ],
377        )? != 1
378        {
379            return Err(DbError::Message(
380                "membership mutation changed during replacement".into(),
381            ));
382        }
383        if let Some(generation) = previous_rotation_generation {
384            super::membership_rotation::remove_rotation_candidate_on(&tx, intent_hash, generation)?;
385        }
386        super::membership_rotation::stage_pending_rotation_on(
387            &tx,
388            pending_rotation_generation,
389            replacement_hash,
390        )?;
391        super::active_store_publication::update_active_store_publication_on(
392            &tx,
393            &expected,
394            &replacement,
395        )?;
396        tx.commit()?;
397        Ok(replacement_hash)
398    }
399
400    fn complete_membership_mutation(&mut self, intent_hash: ObjectHash) -> Result<(), DbError> {
401        let conn = self.conn;
402        let deleted = conn
403            .execute(
404                "DELETE FROM outbound_membership_mutation \
405                 WHERE singleton = 1 AND intent_hash = ?1",
406                [intent_hash.to_string()],
407            )
408            .map_err(DbError::from)?;
409        if deleted != 1 {
410            return Err(DbError::Message(
411                "membership mutation ownership row is absent or changed".to_string(),
412            ));
413        }
414        Ok(())
415    }
416}
417
418impl StoreDatabase {
419    pub async fn stage_membership_candidate_abandonment(
420        &self,
421        intent_hash: ObjectHash,
422        expected: ActiveStorePublication,
423        original: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
424        abandonment: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
425    ) -> Result<ActiveStorePublication, DbError> {
426        self.call_store(move |session| {
427            session.stage_membership_candidate_abandonment(
428                intent_hash,
429                expected,
430                original,
431                abandonment,
432            )
433        })
434        .await
435    }
436
437    pub async fn replace_membership_candidate_mutation(
438        &self,
439        intent_hash: ObjectHash,
440        expected: ActiveStorePublication,
441        candidate: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
442        plan_bytes: Vec<u8>,
443        progress_bytes: Vec<u8>,
444        remote_objects: Vec<coven_protocol::remote_object::ClosedRemoteObject>,
445    ) -> Result<ObjectHash, DbError> {
446        self.call_store(move |session| {
447            session.replace_membership_candidate_mutation(
448                intent_hash,
449                expected,
450                candidate,
451                plan_bytes,
452                progress_bytes,
453                remote_objects,
454            )
455        })
456        .await
457    }
458
459    pub async fn outbound_membership_mutation(
460        &self,
461    ) -> Result<Option<DurableMembershipMutation>, DbError> {
462        self.call_store(|session| session.outbound_membership_mutation())
463            .await
464    }
465
466    pub async fn select_membership_author_stream(
467        &self,
468        author_pubkey: &str,
469        author_owner_grant: &coven_protocol::membership::MembershipGrantId,
470        reusable: std::collections::BTreeSet<coven_protocol::membership::AuthorStreamId>,
471    ) -> Result<coven_protocol::membership::AuthorStreamId, DbError> {
472        self.select_causal_author_stream(
473            format!("membership_author_stream/{author_pubkey}/{author_owner_grant}"),
474            reusable,
475        )
476        .await
477    }
478
479    pub async fn select_causal_author_stream(
480        &self,
481        key: String,
482        reusable: std::collections::BTreeSet<coven_protocol::membership::AuthorStreamId>,
483    ) -> Result<coven_protocol::membership::AuthorStreamId, DbError> {
484        let candidate = coven_protocol::membership::AuthorStreamId::from_digest(
485            ObjectHash::digest(self.new_store_write_id().as_str().as_bytes()),
486        );
487        self.call_store(move |session| {
488            session.select_causal_author_stream(&key, &reusable, candidate)
489        })
490        .await
491    }
492
493    pub async fn stage_membership_candidate_mutation(
494        &self,
495        plan_bytes: Vec<u8>,
496        progress_bytes: Vec<u8>,
497        remote_objects: Vec<coven_protocol::remote_object::ClosedRemoteObject>,
498        candidate: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
499    ) -> Result<ObjectHash, DbError> {
500        self.call_store(move |session| {
501            session.stage_membership_candidate_mutation(
502                plan_bytes,
503                progress_bytes,
504                remote_objects,
505                candidate,
506            )
507        })
508        .await
509    }
510
511    pub async fn update_membership_mutation_progress(
512        &self,
513        intent_hash: ObjectHash,
514        progress_bytes: Vec<u8>,
515    ) -> Result<(), DbError> {
516        self.call_store(move |session| {
517            session.update_membership_mutation_progress(intent_hash, progress_bytes)
518        })
519        .await
520    }
521
522    pub async fn complete_membership_mutation(
523        &self,
524        intent_hash: ObjectHash,
525    ) -> Result<(), DbError> {
526        self.call_store(move |session| session.complete_membership_mutation(intent_hash))
527            .await
528    }
529}
530
531pub(super) fn require_membership_mutation_on(
532    conn: &rusqlite::Connection,
533    intent_hash: ObjectHash,
534) -> Result<(), DbError> {
535    let stored: Vec<u8> = conn.query_row(
536        "SELECT plan_bytes FROM outbound_membership_mutation WHERE singleton = 1 AND intent_hash = ?1",
537        [intent_hash.to_string()],
538        |row| row.get(0),
539    )?;
540    if ObjectHash::digest(&stored) != intent_hash {
541        return Err(DbError::Message(
542            "membership mutation plan differs from its owner".into(),
543        ));
544    }
545    Ok(())
546}
547
548fn validate_membership_candidate_objects(
549    candidate: &coven_protocol::prepared_commit::PreparedStoreOperationCommit,
550    remote_objects: &[coven_protocol::remote_object::ClosedRemoteObject],
551) -> Result<coven_protocol::membership_mutation::PreparedMembershipPublication, DbError> {
552    candidate.validate_closed_shape()?;
553    let publication = candidate.prepared_membership_publication()?;
554    let objects = publication.candidate_object_refs(&candidate.commit, &candidate.reference)?;
555    let supplied = remote_objects
556        .iter()
557        .map(|remote| remote.object().clone())
558        .collect::<BTreeSet<_>>();
559    if supplied.len() != remote_objects.len()
560        || supplied != objects.into_iter().collect::<BTreeSet<_>>()
561    {
562        return Err(DbError::Message(
563            "membership ownership differs from its exact candidate".into(),
564        ));
565    }
566    let owns_candidate = |ownership: &coven_protocol::remote_object::PendingCandidateOwnership| {
567        ownership.pending.len() == 1
568            && ownership.pending.contains(&candidate.reference)
569            && ownership.nonactivated.is_empty()
570    };
571    for remote in remote_objects {
572        use coven_protocol::remote_object::{
573            CandidateCommitState, CandidateObjectState, RetainedAuthorityObjectState,
574        };
575        let prepared_for_candidate = match remote.record() {
576            RemoteObjectRecord::CandidateCommit(record) => {
577                record.identity == candidate.reference
578                    && matches!(record.state, CandidateCommitState::Prepared)
579            }
580            RemoteObjectRecord::CandidateExclusive(record) => matches!(
581                &record.state,
582                CandidateObjectState::Prepared { ownership } if owns_candidate(ownership)
583            ),
584            RemoteObjectRecord::RetainedAuthority(record) => matches!(
585                &record.state,
586                RetainedAuthorityObjectState::Prepared { ownership } if owns_candidate(ownership)
587            ),
588            RemoteObjectRecord::SharedLiveSet(_) => false,
589        };
590        if !prepared_for_candidate {
591            return Err(DbError::Message(
592                "membership object is not prepared for its exact candidate".into(),
593            ));
594        }
595    }
596    Ok(publication)
597}
598
599pub(super) fn membership_rotation_generation(
600    entry: &coven_protocol::membership::MembershipEntry,
601) -> Result<Option<u64>, DbError> {
602    match &entry.change {
603        coven_protocol::membership::StoreAuthorityChange::RemoveMember { wrapped_keys, .. } => {
604            let generation = wrapped_keys
605                .first()
606                .ok_or_else(|| {
607                    DbError::Message(
608                        "retained member removal has no replacement key generation".into(),
609                    )
610                })?
611                .generation;
612            Ok(Some(generation))
613        }
614        _ => Ok(None),
615    }
616}
617
618pub(super) fn require_retired_membership_candidate_on<'a>(
619    conn: &rusqlite::Connection,
620    expected: &'a ActiveStorePublication,
621) -> Result<&'a RetiredStoreCandidate, DbError> {
622    let [retired] = expected.retired_candidates() else {
623        return Err(DbError::Message(
624            "membership continuation has no exact original candidate".into(),
625        ));
626    };
627    if expected.owner() != &ActiveStorePublicationOwner::MembershipMutation
628        || !expected.is_awaiting_preparation()
629        || !matches!(retired.inputs, RetiredStoreCandidateInputs::Membership(_))
630        || super::active_store_publication::load_active_store_publication_on(conn)?.as_ref()
631            != Some(expected)
632    {
633        return Err(DbError::Message(
634            "membership continuation differs from its durable owner".into(),
635        ));
636    }
637    super::candidate_records::require_candidate_cleanup_complete_on(
638        conn,
639        &retired.candidate()?,
640        &retired.objects()?,
641        "membership candidate cleanup is incomplete",
642    )?;
643    Ok(retired)
644}
645
646pub(super) fn consume_retired_membership_candidate_on(
647    tx: &rusqlite::Transaction<'_>,
648    expected: &ActiveStorePublication,
649) -> Result<ActiveStorePublication, DbError> {
650    let retired = require_retired_membership_candidate_on(tx, expected)?;
651    super::candidate_records::delete_remote_objects_on(
652        tx,
653        retired.objects()?.iter().map(remote_object_id),
654        "retired membership candidate",
655    )?;
656    let mut completed = expected.clone();
657    completed.complete_retired_candidate_cleanup()?;
658    Ok(completed)
659}
660
661#[cfg(test)]
662#[path = "membership_mutations_tests.rs"]
663mod tests;