1use crate::*;
2use coven_protocol::store_commit::{
3 StoreAck, StoreAckRef, StoreBatchCommitRef, StoreDeviceRegistrationRef,
4};
5use rusqlite::OptionalExtension;
6
7use super::*;
8use crate::store_ack_records::{load_expected_outbound_store_ack_on, load_outbound_store_ack_on};
9
10impl StoreSession<'_> {
11 fn latest_local_store_ack(&self) -> Result<Option<PublishedStoreAck>, DbError> {
12 load_published_store_ack_on(self.conn)
13 }
14
15 pub(super) fn activated_store_ack(
16 &mut self,
17 registration: &StoreDeviceRegistrationRef,
18 ) -> Result<Option<ActivatedStoreAck>, DbError> {
19 let activated = self
20 .conn
21 .query_row(
22 "SELECT ack_ref, activating_commit FROM activated_store_acks WHERE device_id = ?1",
23 [registration.device_id.to_string()],
24 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
25 )
26 .optional()
27 .map_err(DbError::from)?
28 .map(|(raw, activating_commit)| {
29 let reference: StoreAckRef = serde_json::from_str(&raw).map_err(|error| {
30 DbError::context("activated Store acknowledgement ref", error)
31 })?;
32 if &reference.registration != registration {
33 return Err(DbError::Message(
34 "activated Store acknowledgement names another registration".to_string(),
35 ));
36 }
37 Ok(ActivatedStoreAck {
38 reference,
39 activating_commit: serde_json::from_str(&activating_commit).map_err(
40 |error| {
41 DbError::context(
42 "activated Store acknowledgement activating commit",
43 error,
44 )
45 },
46 )?,
47 })
48 })
49 .transpose()?;
50 let baseline = self
51 .verified_store_authority
52 .retained_replay_baseline_on(StoreRecords::new(self.conn, self.store_dir))?;
53 let RetainedReplayAuthority::InstalledSnapshot(snapshot) = &baseline.authority else {
54 return Ok(activated);
55 };
56 let Some(chain) = snapshot
57 .metadata
58 .history_summary
59 .acknowledgements
60 .get(®istration.device_id)
61 else {
62 return Ok(activated);
63 };
64 let (reference, _) = chain.latest().ok_or_else(|| {
65 DbError::Message("installed snapshot acknowledgement chain is empty".into())
66 })?;
67 if &reference.registration != registration {
68 return Err(DbError::Message(
69 "snapshot acknowledgement names another registration".into(),
70 ));
71 }
72 if let Some(current) = &activated {
73 if current.reference.sequence > reference.sequence {
74 return Ok(activated);
75 }
76 if chain
77 .chain
78 .get(¤t.reference.sequence)
79 .map(|(reference, _)| reference)
80 != Some(¤t.reference)
81 || (current.reference.sequence == reference.sequence
82 && current.activating_commit != chain.activating_commit)
83 {
84 return Err(DbError::Message(
85 "activated acknowledgement conflicts with the installed snapshot".into(),
86 ));
87 }
88 }
89 Ok(Some(ActivatedStoreAck {
90 reference: reference.clone(),
91 activating_commit: chain.activating_commit.clone(),
92 }))
93 }
94
95 fn stage_store_ack(
96 &mut self,
97 ack: StoreAck,
98 prepared: PreparedExactObject,
99 ) -> Result<StoreAckRef, DbError> {
100 let authority = self.local_store_authority()?;
101 let bytes = ack.to_bytes();
102 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
103 let reference = verify_next_local_store_ack_on(&tx, &authority, &bytes, &prepared)?;
104 let ack_ref = serde_json::to_string(&reference).map_err(|error| {
105 DbError::context("serialize exact Store acknowledgement ref", error)
106 })?;
107 let prepared = serde_json::to_string(&prepared)
108 .map_err(|error| DbError::context("serialize prepared Store acknowledgement", error))?;
109 let activation = serde_json::to_string(&OutboundStoreAckActivation::AwaitingCandidate)
110 .map_err(|error| {
111 DbError::context("serialize Store acknowledgement activation state", error)
112 })?;
113 tx.execute(
114 "INSERT INTO outbound_store_acks
115 (singleton, ack_ref, ack_bytes, prepared_object, activation)
116 VALUES (1, ?1, ?2, ?3, ?4)",
117 rusqlite::params![ack_ref, bytes, prepared, activation],
118 )
119 .map_err(DbError::from)?;
120 tx.commit().map_err(DbError::from)?;
121 Ok(reference)
122 }
123
124 fn adopt_outbound_store_ack_slot_winner(
125 &mut self,
126 expected: &StoreAckRef,
127 winner_bytes: Vec<u8>,
128 winner_prepared: PreparedExactObject,
129 ) -> Result<(), DbError> {
130 let authority = self.local_store_authority()?;
131 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
132 let outbound = load_expected_outbound_store_ack_on(
133 &tx,
134 &authority,
135 expected,
136 "acknowledgement slot winner names another queued object",
137 )?;
138 let OutboundStoreAckActivation::Prepared(candidate) = &outbound.activation else {
139 return Err(DbError::Message(
140 "acknowledgement slot collision has no prepared activation candidate".to_string(),
141 ));
142 };
143 let active_publication = ActiveStorePublication::for_commit(
144 crate::ActiveStorePublicationOwner::StoreAcknowledgement,
145 candidate,
146 )?;
147 if candidate.commit.acknowledgement() != Some(expected) {
148 return Err(DbError::Message(
149 "prepared activation candidate names another acknowledgement".to_string(),
150 ));
151 }
152 if winner_prepared.reference().slot() != expected.object.slot()
153 || winner_prepared.reference() == &expected.object
154 {
155 return Err(DbError::Message(
156 "acknowledgement slot winner is not a distinct object at the occupied slot"
157 .to_string(),
158 ));
159 }
160 let winner_reference =
161 verify_next_local_store_ack_on(&tx, &authority, &winner_bytes, &winner_prepared)?;
162 let mut expected_records = candidate
163 .acknowledgement_remote_objects(&outbound.ack)
164 .map_err(DbError::from)?;
165 for reference in candidate.commit.circle_acknowledgements() {
166 let circle = outbound
167 .circle_acknowledgements
168 .iter()
169 .find(|circle| &circle.reference == reference)
170 .ok_or_else(|| {
171 DbError::Message(
172 "losing acknowledgement candidate lost its queued Circle statement".into(),
173 )
174 })?;
175 expected_records.extend(candidate.circle_acknowledgement_remote_objects(&circle.ack)?);
176 }
177 let expected_records = expected_records
178 .into_iter()
179 .map(|record| (record.object_id(), record))
180 .collect::<std::collections::BTreeMap<_, _>>();
181 for expected_record in expected_records.values() {
182 let object_id = expected_record.object_id();
183 let stored = load_remote_object_on(&tx, object_id)?;
184 if stored != **expected_record {
185 return Err(DbError::Message(
186 "losing acknowledgement candidate is no longer wholly unuploaded".to_string(),
187 ));
188 }
189 }
190 for expected_record in expected_records.into_values() {
191 if !crate::remote_object_records::delete_remote_object_on(
192 &tx,
193 expected_record.object_id(),
194 )? {
195 return Err(DbError::Message(
196 "losing acknowledgement candidate object disappeared".to_string(),
197 ));
198 }
199 }
200 let activation =
201 serde_json::to_string(&OutboundStoreAckActivation::Created).map_err(|error| {
202 DbError::context("serialize adopted Store acknowledgement activation", error)
203 })?;
204 let winner_ref = serde_json::to_string(&winner_reference).map_err(|error| {
205 DbError::context("serialize adopted Store acknowledgement ref", error)
206 })?;
207 let winner_prepared = serde_json::to_string(&winner_prepared).map_err(|error| {
208 DbError::context("serialize adopted prepared Store acknowledgement", error)
209 })?;
210 let updated = tx
211 .execute(
212 "UPDATE outbound_store_acks
213 SET ack_ref = ?2, ack_bytes = ?3, prepared_object = ?4, activation = ?5
214 WHERE singleton = 1 AND ack_ref = ?1",
215 rusqlite::params![
216 serde_json::to_string(expected).map_err(|error| {
217 DbError::context("serialize losing Store acknowledgement ref", error)
218 })?,
219 winner_ref,
220 winner_bytes,
221 winner_prepared,
222 activation,
223 ],
224 )
225 .map_err(DbError::from)?;
226 if updated != 1 {
227 return Err(DbError::Message(
228 "outbound Store acknowledgement changed during winner adoption".to_string(),
229 ));
230 }
231 super::active_store_publication::clear_active_store_publication_on(
232 &tx,
233 &active_publication,
234 )?;
235 tx.commit().map_err(DbError::from)
236 }
237
238 fn oldest_outbound_store_ack(&mut self) -> Result<Option<OutboundStoreAck>, DbError> {
239 let authority = self.local_store_authority()?;
240 load_outbound_store_ack_on(self.conn, &authority)
241 }
242
243 fn complete_outbound_store_ack(
244 &mut self,
245 accepted: &StoreAckRef,
246 activating_commit: &StoreBatchCommitRef,
247 ) -> Result<(), DbError> {
248 let authority = self.local_store_authority()?;
249 let activated = self.activated_store_ack(&accepted.registration)?;
250 if !activated.is_some_and(|activated| {
251 &activated.reference == accepted && &activated.activating_commit == activating_commit
252 }) {
253 return Err(DbError::Message(
254 "Store acknowledgement completion requires its exact installed activation"
255 .to_string(),
256 ));
257 }
258 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
259 let outbound = load_expected_outbound_store_ack_on(
260 &tx,
261 &authority,
262 accepted,
263 "accepted Store acknowledgement differs from the prepared exact object",
264 )?;
265 let candidate = match &outbound.activation {
266 OutboundStoreAckActivation::Prepared(candidate) => Some(candidate),
267 OutboundStoreAckActivation::Created => None,
268 OutboundStoreAckActivation::AwaitingCandidate => {
269 return Err(DbError::Message(
270 "accepted Store acknowledgement has no prepared activation".to_string(),
271 ));
272 }
273 };
274 if candidate.is_some_and(|candidate| &candidate.reference != activating_commit) {
275 return Err(DbError::Message(
276 "accepted Store acknowledgement differs from its prepared activating commit"
277 .to_string(),
278 ));
279 }
280 finish_outbound_store_ack_on(
281 &tx,
282 accepted,
283 &outbound.ack.value.successor.next_slot,
284 &coven_protocol::store_commit::StandingStoreAck {
285 assertion: outbound.ack.value.assertion(),
286 activating_commit: Some(activating_commit.clone()),
287 },
288 )?;
289 for circle in outbound.circle_acknowledgements.iter().filter(|circle| {
290 candidate.is_some_and(|candidate| {
291 candidate
292 .commit
293 .circle_acknowledgements()
294 .contains(&circle.reference)
295 })
296 }) {
297 let circle_id = circle.reference.circle_id.to_string();
298 let removed = tx
299 .execute(
300 "DELETE FROM outbound_circle_acks WHERE circle_id = ?1",
301 [&circle_id],
302 )
303 .map_err(DbError::from)?;
304 if removed != 1 {
305 return Err(DbError::Message(
306 "outbound Circle acknowledgement disappeared during completion".to_string(),
307 ));
308 }
309 let successor_slot = serde_json::to_string(&circle.ack.value.successor.next_slot)
310 .map_err(|error| {
311 DbError::context("serialize Circle acknowledgement successor slot", error)
312 })?;
313 let store_cut = serde_json::to_string(&circle.ack.value.store_cut)
314 .map_err(|error| DbError::context("serialize Circle acknowledgement cut", error))?;
315 let control_coord =
316 serde_json::to_string(&circle.ack.value.control).map_err(|error| {
317 DbError::context("serialize Circle acknowledgement control", error)
318 })?;
319 tx.execute(
320 "INSERT INTO published_circle_acks
321 (circle_id, ack_ref, successor_slot, store_cut, control_coord)
322 VALUES (?1, ?2, ?3, ?4, ?5)
323 ON CONFLICT(circle_id) DO UPDATE SET
324 ack_ref = excluded.ack_ref, successor_slot = excluded.successor_slot,
325 store_cut = excluded.store_cut, control_coord = excluded.control_coord",
326 rusqlite::params![
327 circle_id,
328 serde_json::to_string(&circle.reference).map_err(|error| {
329 DbError::context("serialize published Circle acknowledgement ref", error)
330 })?,
331 successor_slot,
332 store_cut,
333 control_coord,
334 ],
335 )
336 .map_err(DbError::from)?;
337 }
338 if candidate.is_some() {
341 super::active_store_publication::clear_active_store_commit_for_owner_on(
342 &tx,
343 &crate::ActiveStorePublicationOwner::StoreAcknowledgement,
344 activating_commit,
345 )?;
346 }
347 tx.commit().map_err(DbError::from)
348 }
349}
350
351impl StoreDatabase {
352 pub async fn latest_local_store_ack(&self) -> Result<Option<PublishedStoreAck>, DbError> {
353 self.call_store(|session| session.latest_local_store_ack())
354 .await
355 }
356
357 pub async fn activated_store_ack(
358 &self,
359 registration: &StoreDeviceRegistrationRef,
360 ) -> Result<Option<ActivatedStoreAck>, DbError> {
361 let registration = registration.clone();
362 self.call_store(move |session| session.activated_store_ack(®istration))
363 .await
364 }
365
366 pub async fn stage_store_ack(
367 &self,
368 ack: StoreAck,
369 prepared: PreparedExactObject,
370 ) -> Result<StoreAckRef, DbError> {
371 self.call_store(move |session| session.stage_store_ack(ack, prepared))
372 .await
373 }
374
375 pub async fn adopt_outbound_store_ack_slot_winner(
376 &self,
377 expected: StoreAckRef,
378 winner_bytes: Vec<u8>,
379 winner_prepared: PreparedExactObject,
380 ) -> Result<(), DbError> {
381 self.call_store(move |session| {
382 session.adopt_outbound_store_ack_slot_winner(&expected, winner_bytes, winner_prepared)
383 })
384 .await
385 }
386
387 pub async fn oldest_outbound_store_ack(&self) -> Result<Option<OutboundStoreAck>, DbError> {
388 self.call_store(|session| session.oldest_outbound_store_ack())
389 .await
390 }
391
392 pub async fn complete_outbound_store_ack(
393 &self,
394 accepted: StoreAckRef,
395 activating_commit: StoreBatchCommitRef,
396 ) -> Result<(), DbError> {
397 self.call_store(move |session| {
398 session.complete_outbound_store_ack(&accepted, &activating_commit)
399 })
400 .await
401 }
402}