Skip to main content

coven_database/store/store_session/
acknowledgements.rs

1use 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 (&current, 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(), &current)?;
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}