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;