coven_database/store/store_session/
store_acknowledgements.rs1use 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 fn activated_store_ack(
16 &self,
17 registration: &StoreDeviceRegistrationRef,
18 ) -> Result<Option<ActivatedStoreAck>, DbError> {
19 self.conn
20 .query_row(
21 "SELECT ack_ref, activating_commit FROM activated_store_acks WHERE device_id = ?1",
22 [registration.device_id.to_string()],
23 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
24 )
25 .optional()
26 .map_err(DbError::from)?
27 .map(|(raw, activating_commit)| {
28 let reference: StoreAckRef = serde_json::from_str(&raw).map_err(|error| {
29 DbError::context("activated Store acknowledgement ref", error)
30 })?;
31 if &reference.registration != registration {
32 return Err(DbError::Message(
33 "activated Store acknowledgement names another registration".to_string(),
34 ));
35 }
36 Ok(ActivatedStoreAck {
37 reference,
38 activating_commit: serde_json::from_str(&activating_commit).map_err(
39 |error| {
40 DbError::context(
41 "activated Store acknowledgement activating commit",
42 error,
43 )
44 },
45 )?,
46 })
47 })
48 .transpose()
49 }
50
51 fn stage_store_ack(
52 &mut self,
53 ack: StoreAck,
54 prepared: PreparedExactObject,
55 ) -> Result<StoreAckRef, DbError> {
56 let authority = self.local_store_authority()?;
57 let bytes = ack.to_bytes();
58 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
59 let (reference, verified) =
60 verify_next_local_store_ack_on(&tx, &authority, &bytes, &prepared)?;
61 if verified != ack {
62 return Err(DbError::Message(
63 "staged Store acknowledgement changed during exact verification".to_string(),
64 ));
65 }
66 let ack_ref = serde_json::to_string(&reference).map_err(|error| {
67 DbError::context("serialize exact Store acknowledgement ref", error)
68 })?;
69 let prepared = serde_json::to_string(&prepared)
70 .map_err(|error| DbError::context("serialize prepared Store acknowledgement", error))?;
71 let activation = serde_json::to_string(&OutboundStoreAckActivation::AwaitingCandidate)
72 .map_err(|error| {
73 DbError::context("serialize Store acknowledgement activation state", error)
74 })?;
75 tx.execute(
76 "INSERT INTO outbound_store_acks
77 (singleton, ack_ref, ack_bytes, prepared_object, activation)
78 VALUES (1, ?1, ?2, ?3, ?4)",
79 rusqlite::params![ack_ref, bytes, prepared, activation],
80 )
81 .map_err(DbError::from)?;
82 tx.commit().map_err(DbError::from)?;
83 Ok(reference)
84 }
85
86 fn adopt_outbound_store_ack_slot_winner(
87 &mut self,
88 expected: &StoreAckRef,
89 winner_bytes: Vec<u8>,
90 winner_prepared: PreparedExactObject,
91 ) -> Result<(), DbError> {
92 let authority = self.local_store_authority()?;
93 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
94 let outbound = load_expected_outbound_store_ack_on(
95 &tx,
96 &authority,
97 expected,
98 "acknowledgement slot winner names another queued object",
99 )?;
100 let OutboundStoreAckActivation::Prepared(candidate) = &outbound.activation else {
101 return Err(DbError::Message(
102 "acknowledgement slot collision has no prepared activation candidate".to_string(),
103 ));
104 };
105 if candidate.commit.acknowledgement() != Some(expected) {
106 return Err(DbError::Message(
107 "prepared activation candidate names another acknowledgement".to_string(),
108 ));
109 }
110 if winner_prepared.reference().slot() != expected.object.slot()
111 || winner_prepared.reference() == &expected.object
112 {
113 return Err(DbError::Message(
114 "acknowledgement slot winner is not a distinct object at the occupied slot"
115 .to_string(),
116 ));
117 }
118 let (winner_reference, _) =
119 verify_next_local_store_ack_on(&tx, &authority, &winner_bytes, &winner_prepared)?;
120 let expected_records = candidate
121 .acknowledgement_remote_objects(&outbound.ack)
122 .map_err(DbError::from)?;
123 for expected_record in &expected_records {
124 let object_id = expected_record.object_id();
125 let stored = load_remote_object_on(&tx, object_id)?;
126 if stored != **expected_record {
127 return Err(DbError::Message(
128 "losing acknowledgement candidate is no longer wholly unuploaded".to_string(),
129 ));
130 }
131 }
132 for expected_record in expected_records {
133 if !crate::remote_object_records::delete_remote_object_on(
134 &tx,
135 expected_record.object_id(),
136 )? {
137 return Err(DbError::Message(
138 "losing acknowledgement candidate object disappeared".to_string(),
139 ));
140 }
141 }
142 let activation = serde_json::to_string(&OutboundStoreAckActivation::AwaitingCandidate)
143 .map_err(|error| {
144 DbError::context("serialize adopted Store acknowledgement activation", error)
145 })?;
146 let winner_ref = serde_json::to_string(&winner_reference).map_err(|error| {
147 DbError::context("serialize adopted Store acknowledgement ref", error)
148 })?;
149 let winner_prepared = serde_json::to_string(&winner_prepared).map_err(|error| {
150 DbError::context("serialize adopted prepared Store acknowledgement", error)
151 })?;
152 let updated = tx
153 .execute(
154 "UPDATE outbound_store_acks
155 SET ack_ref = ?2, ack_bytes = ?3, prepared_object = ?4, activation = ?5
156 WHERE singleton = 1 AND ack_ref = ?1",
157 rusqlite::params![
158 serde_json::to_string(expected).map_err(|error| {
159 DbError::context("serialize losing Store acknowledgement ref", error)
160 })?,
161 winner_ref,
162 winner_bytes,
163 winner_prepared,
164 activation,
165 ],
166 )
167 .map_err(DbError::from)?;
168 if updated != 1 {
169 return Err(DbError::Message(
170 "outbound Store acknowledgement changed during winner adoption".to_string(),
171 ));
172 }
173 tx.commit().map_err(DbError::from)
174 }
175
176 fn oldest_outbound_store_ack(&mut self) -> Result<Option<OutboundStoreAck>, DbError> {
177 let authority = self.local_store_authority()?;
178 load_outbound_store_ack_on(self.conn, &authority)
179 }
180
181 fn complete_outbound_store_ack(
182 &mut self,
183 accepted: &StoreAckRef,
184 activating_commit: &StoreBatchCommitRef,
185 ) -> Result<(), DbError> {
186 let authority = self.local_store_authority()?;
187 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
188 let outbound = load_expected_outbound_store_ack_on(
189 &tx,
190 &authority,
191 accepted,
192 "accepted Store acknowledgement differs from the prepared exact object",
193 )?;
194 finish_outbound_store_ack_on(
195 &tx,
196 accepted,
197 &outbound.ack.value.successor.next_slot,
198 &coven_protocol::store_commit::StandingStoreAck {
199 assertion: outbound.ack.value.assertion(),
200 activating_commit: Some(activating_commit.clone()),
201 },
202 )?;
203 for circle in &outbound.circle_acknowledgements {
204 let circle_id = circle.reference.circle_id.to_string();
205 let removed = tx
206 .execute(
207 "DELETE FROM outbound_circle_acks WHERE circle_id = ?1",
208 [&circle_id],
209 )
210 .map_err(DbError::from)?;
211 if removed != 1 {
212 return Err(DbError::Message(
213 "outbound Circle acknowledgement disappeared during completion".to_string(),
214 ));
215 }
216 let successor_slot = serde_json::to_string(&circle.ack.value.successor.next_slot)
217 .map_err(|error| {
218 DbError::context("serialize Circle acknowledgement successor slot", error)
219 })?;
220 let store_cut = serde_json::to_string(&circle.ack.value.store_cut)
221 .map_err(|error| DbError::context("serialize Circle acknowledgement cut", error))?;
222 let control_coord =
223 serde_json::to_string(&circle.ack.value.control).map_err(|error| {
224 DbError::context("serialize Circle acknowledgement control", error)
225 })?;
226 tx.execute(
227 "INSERT INTO published_circle_acks
228 (circle_id, ack_ref, successor_slot, store_cut, control_coord)
229 VALUES (?1, ?2, ?3, ?4, ?5)
230 ON CONFLICT(circle_id) DO UPDATE SET
231 ack_ref = excluded.ack_ref, successor_slot = excluded.successor_slot,
232 store_cut = excluded.store_cut, control_coord = excluded.control_coord",
233 rusqlite::params![
234 circle_id,
235 serde_json::to_string(&circle.reference).map_err(|error| {
236 DbError::context("serialize published Circle acknowledgement ref", error)
237 })?,
238 successor_slot,
239 store_cut,
240 control_coord,
241 ],
242 )
243 .map_err(DbError::from)?;
244 }
245 tx.commit().map_err(DbError::from)
246 }
247}
248
249impl StoreDatabase {
250 pub async fn latest_local_store_ack(&self) -> Result<Option<PublishedStoreAck>, DbError> {
251 self.call_store(|session| session.latest_local_store_ack())
252 .await
253 }
254
255 pub async fn activated_store_ack(
256 &self,
257 registration: &StoreDeviceRegistrationRef,
258 ) -> Result<Option<ActivatedStoreAck>, DbError> {
259 let registration = registration.clone();
260 self.call_store(move |session| session.activated_store_ack(®istration))
261 .await
262 }
263
264 pub async fn stage_store_ack(
265 &self,
266 ack: StoreAck,
267 prepared: PreparedExactObject,
268 ) -> Result<StoreAckRef, DbError> {
269 self.call_store(move |session| session.stage_store_ack(ack, prepared))
270 .await
271 }
272
273 pub async fn adopt_outbound_store_ack_slot_winner(
274 &self,
275 expected: StoreAckRef,
276 winner_bytes: Vec<u8>,
277 winner_prepared: PreparedExactObject,
278 ) -> Result<(), DbError> {
279 self.call_store(move |session| {
280 session.adopt_outbound_store_ack_slot_winner(&expected, winner_bytes, winner_prepared)
281 })
282 .await
283 }
284
285 pub async fn oldest_outbound_store_ack(&self) -> Result<Option<OutboundStoreAck>, DbError> {
286 self.call_store(|session| session.oldest_outbound_store_ack())
287 .await
288 }
289
290 pub async fn complete_outbound_store_ack(
291 &self,
292 accepted: StoreAckRef,
293 activating_commit: StoreBatchCommitRef,
294 ) -> Result<(), DbError> {
295 self.call_store(move |session| {
296 session.complete_outbound_store_ack(&accepted, &activating_commit)
297 })
298 .await
299 }
300}