coven_database/store/store_session/
acknowledgements.rs1use super::*;
2use crate::store_ack_records::{
3 load_expected_outbound_store_ack_on, verify_next_local_store_ack_on,
4};
5use crate::{
6 update_remote_object_on, ActiveStorePublication, ActiveStorePublicationOwner,
7 ExactProtocolObject, RemoteObjectRecord, StoreAck,
8};
9
10impl StoreSession<'_> {
11 fn replace_acknowledgement_activation(
12 &mut self,
13 expected: StoreAckRef,
14 snapshot: coven_protocol::store_commit::AcceptedStoreSnapshotRef,
15 acknowledgement: ExactProtocolObject<StoreAck>,
16 candidate: PreparedStoreOperationCommit,
17 ) -> Result<(), DbError> {
18 let authority = self.local_store_authority()?;
19 self.verified_store_transaction(move |transaction| {
20 let tx = transaction.store.transaction;
21 let outbound = load_expected_outbound_store_ack_on(
22 tx, &authority, &expected,
23 "acknowledgement replacement names another queued object",
24 )?;
25 let OutboundStoreAckActivation::Prepared(previous) = &outbound.activation else {
26 return Err(DbError::Message("acknowledgement replacement has no prepared candidate".into()));
27 };
28 previous.validate_closed_shape()?;
29 candidate.validate_closed_shape()?;
30 let old_commit = coven_protocol::store_commit::VerifiedStoreBatchCommit::parse_prepared(
31 &previous.commit.to_bytes(), authority.value().store_root.store_root_hash,
32 previous.reference.coord.clone(), previous.reference.object.clone(), authority.value(),
33 )?;
34 let new_commit = coven_protocol::store_commit::VerifiedStoreBatchCommit::parse_prepared(
35 &candidate.commit.to_bytes(), authority.value().store_root.store_root_hash,
36 candidate.reference.coord.clone(), candidate.reference.object.clone(), authority.value(),
37 )?;
38 let proof = transaction.snapshot_candidate_nonactivation(&snapshot, &old_commit)?;
39 let active = super::active_store_publication::load_active_store_publication_on(tx)?
40 .ok_or_else(|| DbError::Message("acknowledgement replacement has no reserved publication".into()))?;
41 let installed = super::observed_store_publication::load_store_current_publication_on(tx)?;
42 if installed.record() != &candidate.publication.previous
43 || installed.observed_version() != Some(&candidate.publication.previous_version)
44 || new_commit.publication_base != coven_protocol::store_commit::StorePublicationBase::Snapshot(snapshot)
45 {
46 return Err(DbError::Message("replacement acknowledgement extends another installed boundary".into()));
47 }
48 candidate.publication.verify_commit(&new_commit)?;
49 let retained = candidate.history_evidence.acknowledgement.as_ref()
50 .ok_or_else(|| DbError::Message("replacement acknowledgement omits its proof".into()))?;
51 let previous_proof = previous.history_evidence.acknowledgement.as_ref()
52 .ok_or_else(|| DbError::Message("queued acknowledgement omits its proof".into()))?;
53 if previous_proof.acknowledgement != (outbound.reference.clone(), outbound.ack.value.clone())
54 || retained.predecessors != previous_proof.proof_objects().cloned().collect::<Vec<_>>()
55 || retained.acknowledgement.1 != acknowledgement.value
56 || retained.acknowledgement.0.object != *acknowledgement.prepared.reference()
57 || acknowledgement.value.to_bytes() != acknowledgement.bytes
58 || acknowledgement.value.store_cut != new_commit.order.predecessor_cut()?
59 || acknowledgement.value.device_state != new_commit.device_state
60 || candidate.commit.circle_acknowledgements() != previous.commit.circle_acknowledgements()
61 {
62 return Err(DbError::Message("replacement acknowledgement changes its queued proof or Circle statements".into()));
63 }
64 for (reference, value) in retained.proof_objects() {
65 StoreAck::parse_at(&value.to_bytes(), &authority.value().store_root, reference, authority.value())?;
66 }
67 let mut publications = vec![active.attempt()?.reference()?];
68 if let Some(prior) = active.superseded_entry() {
69 publications.push(prior.clone());
70 }
71 let cleanup = crate::RetiredStoreCandidate {
72 nonactivation: proof,
73 inputs: crate::RetiredStoreCandidateInputs::Acknowledgement((**previous_proof).clone()),
74 publications,
75 };
76 let replacement = active.replace_acknowledgement_candidate(&candidate, cleanup.clone())?;
77 let old_objects = cleanup.objects()?.into_iter().collect::<std::collections::BTreeSet<_>>();
78 let mut objects = candidate.acknowledgement_remote_objects(&acknowledgement)?;
79 for circle in &outbound.circle_acknowledgements {
80 objects.extend(candidate.circle_acknowledgement_remote_objects(&circle.ack)?);
81 }
82 let mut persisted = std::collections::BTreeSet::new();
83 for proposed in objects {
84 if !persisted.insert(proposed.object_id()) {
85 continue;
86 }
87 if old_objects.contains(proposed.object()) {
88 let mut current = load_remote_object_on(tx, proposed.object_id())?;
89 match (¤t, proposed.record()) {
90 (RemoteObjectRecord::RetainedAuthority(current), RemoteObjectRecord::RetainedAuthority(proposed))
91 if current.identity == proposed.identity && current.payloads == proposed.payloads => {}
92 _ => return Err(DbError::Message("replacement changed a retained acknowledgement object".into())),
93 }
94 current.add_retained_authority_candidate(candidate.reference.clone())?;
95 update_remote_object_on(tx, proposed.object_id(), ¤t)?;
96 } else {
97 persist_exact_remote_object_on(tx, transaction.store.store_dir, &proposed, "replacement acknowledgement candidate")?;
98 }
99 }
100 super::candidate_records::begin_candidate_nonactivation_targets_on(
101 tx, &previous.reference, &cleanup.objects()?, &cleanup.nonactivation,
102 )?;
103 let changed = tx.execute(
104 "UPDATE outbound_store_acks SET ack_ref = ?2, ack_bytes = ?3, prepared_object = ?4, activation = ?5 WHERE singleton = 1 AND ack_ref = ?1",
105 rusqlite::params![
106 serde_json::to_string(&expected)?,
107 serde_json::to_string(&retained.acknowledgement.0)?,
108 acknowledgement.bytes,
109 serde_json::to_string(&acknowledgement.prepared)?,
110 serde_json::to_string(&OutboundStoreAckActivation::Prepared(candidate))?,
111 ],
112 )?;
113 if changed != 1 {
114 return Err(DbError::Message("outbound acknowledgement changed during candidate replacement".into()));
115 }
116 super::active_store_publication::update_active_store_publication_on(tx, &active, &replacement)?;
117 Ok(StoreTransactionOutcome::Commit(()))
118 })
119 }
120
121 fn prepare_acknowledgement_activation(
122 &mut self,
123 expected: &StoreAckRef,
124 acknowledgement: ExactProtocolObject<StoreAck>,
125 candidate: PreparedStoreOperationCommit,
126 ) -> Result<bool, DbError> {
127 let authority = self.local_store_authority()?;
128 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
129 let outbound = load_expected_outbound_store_ack_on(
130 &tx,
131 &authority,
132 expected,
133 "prepared activation names a different Store acknowledgement",
134 )?;
135 match &outbound.activation {
136 OutboundStoreAckActivation::AwaitingCandidate | OutboundStoreAckActivation::Created => {
137 }
138 OutboundStoreAckActivation::Prepared(existing)
139 if *existing == candidate
140 && outbound.ack.bytes == acknowledgement.bytes
141 && outbound.ack.prepared == acknowledgement.prepared
142 && outbound.ack.value == acknowledgement.value =>
143 {
144 return Ok(true);
145 }
146 OutboundStoreAckActivation::Prepared(_) => {
147 return Err(DbError::Message(
148 "Store acknowledgement already has a different activation candidate"
149 .to_string(),
150 ));
151 }
152 }
153 candidate.validate_closed_shape()?;
154 let reference = candidate.commit.acknowledgement().ok_or_else(|| {
155 DbError::Message("prepared acknowledgement has no exact statement".into())
156 })?;
157 let value = StoreAck::parse_at(
158 &acknowledgement.bytes,
159 &authority.value().store_root,
160 reference,
161 authority.value(),
162 )?;
163 let retained = candidate
164 .history_evidence
165 .acknowledgement
166 .as_ref()
167 .ok_or_else(|| {
168 DbError::Message("prepared acknowledgement omits its retained statement".into())
169 })?;
170 if value != acknowledgement.value
171 || reference.object != *acknowledgement.prepared.reference()
172 || retained.acknowledgement != (reference.clone(), value.clone())
173 || value.store_cut != candidate.commit.order.predecessor_cut()?
174 || value.device_state != candidate.commit.device_state
175 || value.last_sync != outbound.ack.value.last_sync
176 || candidate.commit.circle_acknowledgements()
177 != outbound
178 .circle_acknowledgements
179 .iter()
180 .map(|circle| circle.reference.clone())
181 .collect::<Vec<_>>()
182 {
183 return Err(DbError::Message("prepared acknowledgement differs from its queued statement or accepted predecessor".into()));
184 }
185 let created = match &outbound.activation {
186 OutboundStoreAckActivation::AwaitingCandidate => {
187 let verified = verify_next_local_store_ack_on(
188 &tx,
189 &authority,
190 &acknowledgement.bytes,
191 &acknowledgement.prepared,
192 )?;
193 if verified != *reference
194 || reference.object.slot() != expected.object.slot()
195 || value.successor != outbound.ack.value.successor
196 || !retained.predecessors.is_empty()
197 {
198 return Err(DbError::Message(
199 "uncreated acknowledgement changed its reserved stream position".into(),
200 ));
201 }
202 None
203 }
204 OutboundStoreAckActivation::Created => {
205 if reference == expected {
206 if acknowledgement.bytes != outbound.ack.bytes
207 || acknowledgement.prepared != outbound.ack.prepared
208 || acknowledgement.value != outbound.ack.value
209 || !retained.predecessors.is_empty()
210 {
211 return Err(DbError::Message(
212 "created acknowledgement bytes cannot change".into(),
213 ));
214 }
215 } else if reference.sequence
216 != expected.sequence.checked_add(1).ok_or_else(|| {
217 DbError::Message("acknowledgement sequence overflow".into())
218 })?
219 || reference.object.slot() != &outbound.ack.value.successor.next_slot
220 || value.successor.predecessor.as_ref() != Some(&expected.object)
221 || retained.predecessors != vec![(expected.clone(), outbound.ack.value.clone())]
222 {
223 return Err(DbError::Message(
224 "created acknowledgement replacement omits its exact predecessor".into(),
225 ));
226 }
227 Some(&expected.object)
228 }
229 OutboundStoreAckActivation::Prepared(_) => {
230 unreachable!("prepared candidates returned above")
231 }
232 };
233 let active_publication = ActiveStorePublication::for_commit(
234 ActiveStorePublicationOwner::StoreAcknowledgement,
235 &candidate,
236 )?;
237 match super::active_store_publication::claim_active_store_publication_on(
238 &tx,
239 &active_publication,
240 )? {
241 super::active_store_publication::ActiveStorePublicationClaim::Acquired => {}
242 super::active_store_publication::ActiveStorePublicationClaim::AlreadyOwned => {
243 return Err(DbError::Message(
244 "Store acknowledgement candidate already owns publication before its journal"
245 .to_string(),
246 ));
247 }
248 super::active_store_publication::ActiveStorePublicationClaim::Occupied(_) => {
249 return Ok(false);
250 }
251 }
252 for remote in candidate
253 .acknowledgement_remote_objects(&acknowledgement)
254 .map_err(DbError::from)?
255 {
256 persist_exact_remote_object_on(
257 &tx,
258 self.store_dir,
259 &remote,
260 "Merge Store acknowledgement activation object",
261 )?;
262 if created == Some(remote.object()) {
263 let object_id = remote.object_id();
264 let mut uploaded = remote.into_record();
265 uploaded.mark_uploaded_verified()?;
266 update_remote_object_on(&tx, object_id, &uploaded)?;
267 }
268 }
269 for circle in &outbound.circle_acknowledgements {
270 for remote in candidate
271 .circle_acknowledgement_remote_objects(&circle.ack)
272 .map_err(DbError::from)?
273 {
274 persist_exact_remote_object_on(
275 &tx,
276 self.store_dir,
277 &remote,
278 "Merge Circle acknowledgement activation object",
279 )?;
280 }
281 }
282 let changed = tx.execute(
283 "UPDATE outbound_store_acks SET ack_ref = ?2, ack_bytes = ?3, prepared_object = ?4, activation = ?5 WHERE singleton = 1 AND ack_ref = ?1",
284 rusqlite::params![
285 serde_json::to_string(expected)?,
286 serde_json::to_string(reference)?,
287 acknowledgement.bytes,
288 serde_json::to_string(&acknowledgement.prepared)?,
289 serde_json::to_string(&OutboundStoreAckActivation::Prepared(candidate))?,
290 ],
291 )?;
292 if changed != 1 {
293 return Err(DbError::Message(
294 "outbound acknowledgement changed during activation preparation".into(),
295 ));
296 }
297 tx.commit().map_err(DbError::from)?;
298 Ok(true)
299 }
300}
301impl StoreDatabase {
302 pub async fn replace_acknowledgement_activation(
303 &self,
304 expected: StoreAckRef,
305 snapshot: coven_protocol::store_commit::AcceptedStoreSnapshotRef,
306 acknowledgement: ExactProtocolObject<StoreAck>,
307 candidate: PreparedStoreOperationCommit,
308 ) -> Result<(), DbError> {
309 self.call_store(move |session| {
310 session.replace_acknowledgement_activation(
311 expected,
312 snapshot,
313 acknowledgement,
314 candidate,
315 )
316 })
317 .await
318 }
319
320 pub async fn prepare_acknowledgement_activation(
321 &self,
322 expected: StoreAckRef,
323 acknowledgement: ExactProtocolObject<StoreAck>,
324 candidate: PreparedStoreOperationCommit,
325 ) -> Result<bool, DbError> {
326 self.call_store(move |session| {
327 session.prepare_acknowledgement_activation(&expected, acknowledgement, candidate)
328 })
329 .await
330 }
331}