1use super::{
2 publication_state::{PreparedStoreWriteState, StoreWritePreparation},
3 StoreDatabase, StoreSession,
4};
5#[cfg(any(test, feature = "test-utils"))]
6use crate::StoreWriteBase;
7use crate::{
8 persist_exact_remote_object_on, ActiveStorePublication, ActiveStorePublicationOwner, DbError,
9 DurablePreparedProtocolObject, LOCAL_DEVICE_ID_STATE_KEY,
10};
11use coven_protocol::remote_object::RemoteObjectRecord;
12use coven_protocol::store_commit::{CommitFrontier, StoreCommitCoord, StoreDeviceRegistrationRef};
13use coven_protocol::write::WriteStatus;
14
15impl StoreSession<'_> {
16 fn table_schema_for_apply(&mut self) -> Result<crate::TableSchema, DbError> {
17 crate::TableSchema::for_apply(self.conn, self.synced_tables, self.gates)
18 }
19
20 fn prepare_store_write_commit(&mut self, stage: StoreWritePreparation) -> Result<(), DbError> {
21 let author = self.activated_registration(&stage.commit.value.author_registration)?;
22 if author.value().store_root != stage.root {
23 return Err(DbError::Message(
24 "prepared Store write belongs to another verified Store root".to_string(),
25 ));
26 }
27 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
28 let local_device_id = crate::required_protocol_state_on(&tx, LOCAL_DEVICE_ID_STATE_KEY)?;
29 let registration_object: String = tx
30 .query_row(
31 "SELECT registration_object \
32 FROM store_device_registration_activations WHERE device_id = ?1",
33 [&local_device_id],
34 |row| row.get(0),
35 )
36 .map_err(DbError::from)?;
37 let registration_ref: StoreDeviceRegistrationRef =
38 serde_json::from_str(®istration_object).map_err(|error| {
39 DbError::context("prepared write exact registration ref", error)
40 })?;
41 if registration_ref != stage.commit.value.author_registration {
42 return Err(DbError::Message(
43 "prepared Store commit author registration differs from local activation"
44 .to_string(),
45 ));
46 }
47 let registration = stage.commit.value.author();
48 if author.value() != registration {
49 return Err(DbError::Message(
50 "prepared write author registration differs from its activated bytes".to_string(),
51 ));
52 }
53 let stream_id = coven_protocol::store_commit::StreamActivation::device_authorized_stream_id(
54 stage.root.store_root_hash,
55 ®istration_ref,
56 coven_protocol::store_commit::StreamAnchorDomain::StoreAnnouncements,
57 );
58 let expected_coord = StoreCommitCoord {
59 stream_id,
60 sequence: stage.commit.value.seq(),
61 };
62 if stage.commit.value.store_root_hash() != stage.root.store_root_hash
63 || stage.commit.value.reference().coord != expected_coord
64 || stage.commit.value.reference().object != *stage.commit.prepared.reference()
65 {
66 return Err(DbError::Message(
67 "authenticated prepared Store commit differs from its current local authority"
68 .to_string(),
69 ));
70 }
71 if stage.commit.value.write_id != stage.write_id {
72 return Err(DbError::Message(
73 "prepared write id differs from signed commit".to_string(),
74 ));
75 }
76 let commit_ref = stage.commit.value.reference().clone();
77 stage
78 .history_evidence
79 .validate_for(&commit_ref, stage.commit.value.value())
80 .map_err(|error| DbError::context("prepared Store history evidence", error))?;
81 let installed_publication =
82 super::observed_store_publication::load_store_current_publication_on(&tx)?;
83 if installed_publication.record() != &stage.publication.previous
84 || installed_publication.observed_version() != Some(&stage.publication.previous_version)
85 {
86 return Err(DbError::Message(
87 "prepared Store write extends a stale publication boundary".to_string(),
88 ));
89 }
90 let publication_reference = coven_protocol::store_commit::StorePublicationRef::from_entry(
91 &stage.publication.entry,
92 stage.publication.entry_object.clone(),
93 )
94 .map_err(|error| DbError::context("prepared Store publication entry", error))?;
95 stage
96 .publication
97 .replacement
98 .verify_commit_transition(
99 &stage.publication.previous,
100 &stage.publication.entry,
101 &publication_reference,
102 &stage.commit.value,
103 ®istration.device_signing_pubkey,
104 )
105 .map_err(|error| DbError::context("prepared Store publication transition", error))?;
106 let (stored_base, stored_status, stored_preparation): (
107 Option<String>,
108 String,
109 Option<String>,
110 ) = tx
111 .query_row(
112 "SELECT base, status, prepared
113 FROM store_writes WHERE write_id = ?1",
114 [stage.write_id.as_str()],
115 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
116 )
117 .map_err(DbError::from)?;
118 if stored_status != "\"pending\"" || stored_preparation.is_some() {
119 return Err(DbError::Message(format!(
120 "write {} is not an unprepared pending write",
121 stage.write_id
122 )));
123 }
124 let stored_base = stored_base.ok_or_else(|| {
125 DbError::Message(format!(
126 "pending write {} carries no commit base",
127 stage.write_id
128 ))
129 })?;
130 let partitions = crate::store::store_session::StoreTransaction::new(&tx, self.store_dir)
131 .store_write_partitions(stage.write_id.as_str())?;
132 let records = super::StoreRecords::new(&tx, self.store_dir);
133 let stored_base = records.effective_store_write_base(&stage.write_id, &stored_base)?;
134 if let Some(rebased) = records.rebased_store_write(&stage.write_id)? {
135 if rebased.publication_base != stage.commit.value.publication_base {
136 return Err(DbError::Message(
137 "rebased write requires the validated snapshot base".to_string(),
138 ));
139 }
140 }
141 let mut stored_dependencies = CommitFrontier::from_refs(stored_base.dependencies)
142 .map_err(|error| DbError::context("stored write dependencies", error))?;
143 let observed_predecessor = stored_dependencies.0.remove(&stream_id);
144 if stored_dependencies.commits() != stage.commit.value.merge_dependencies() {
145 return Err(DbError::Message(format!(
146 "prepared commit dependencies differ from write {}",
147 stage.write_id
148 )));
149 }
150 let durable_predecessor =
151 crate::store::materialized_commit_index::latest_position_for_device_on(
152 &tx,
153 &stream_id.to_string(),
154 )?;
155 if observed_predecessor.as_ref().is_some_and(|observed| {
156 durable_predecessor.as_ref().is_none_or(|current| {
157 current.coord.sequence() < observed.coord.sequence()
158 || current.coord.sequence() == observed.coord.sequence() && current != observed
159 })
160 }) {
161 return Err(DbError::Message(format!(
162 "outbound Store commit predecessor does not cover write {} capture frontier",
163 stage.write_id
164 )));
165 }
166 let expected_seq = durable_predecessor
167 .as_ref()
168 .map_or(1, |reference| reference.coord.sequence().saturating_add(1));
169 if stage.commit.value.seq() != expected_seq
170 || stage.commit.value.order.predecessor() != durable_predecessor.as_ref()
171 {
172 return Err(DbError::Message(format!(
173 "outbound Store commit exact predecessor differs from durable {durable_predecessor:?}"
174 )));
175 }
176
177 let active_publication = ActiveStorePublication::commit(
178 ActiveStorePublicationOwner::StoreWrite(stage.write_id.clone()),
179 stage.write_id.clone(),
180 registration_ref,
181 commit_ref.coord.clone(),
182 stage.publication.clone(),
183 )?;
184 let reserved = super::active_store_publication::load_active_store_publication_on(&tx)?;
185 if let Some(reserved) = reserved
186 .as_ref()
187 .filter(|active| active.is_awaiting_preparation())
188 {
189 if reserved.owner() != active_publication.owner()
190 || reserved.commit_reservation() != active_publication.commit_reservation()
191 {
192 return Err(DbError::Message(
193 "candidate preparation differs from its retained reservation".to_string(),
194 ));
195 }
196 let replacement = reserved.replace_attempt(stage.publication.clone())?;
197 super::active_store_publication::update_active_store_publication_on(
198 &tx,
199 reserved,
200 &replacement,
201 )?;
202 } else {
203 match super::active_store_publication::claim_active_store_publication_on(
204 &tx,
205 &active_publication,
206 )? {
207 super::active_store_publication::ActiveStorePublicationClaim::Acquired => {}
208 super::active_store_publication::ActiveStorePublicationClaim::AlreadyOwned => {
209 return Err(DbError::Message(
210 "pending Store write already owns an active publication".to_string(),
211 ));
212 }
213 super::active_store_publication::ActiveStorePublicationClaim::Occupied(source) => {
214 return Err(DbError::Message(format!(
215 "another local Store operation owns publication: {source:?}"
216 )));
217 }
218 }
219 }
220
221 let mut object_ids = std::collections::BTreeSet::new();
222 for remote in &stage.remote_objects {
223 remote
224 .validate()
225 .map_err(|error| DbError::context("prepared remote object", error))?;
226 if !object_ids.insert(remote.object_id()) {
227 return Err(DbError::Message(
228 "prepared write contains a duplicate remote object".to_string(),
229 ));
230 }
231 }
232 crate::validate_prepared_audience_blob_graph(&object_ids, &stage.audiences)?;
233 for remote in &stage.remote_objects {
234 crate::persist_prepared_remote_object_on(
235 &tx,
236 self.store_dir,
237 remote,
238 &commit_ref,
239 "candidate audience object",
240 )?;
241 }
242 let commit_remote = RemoteObjectRecord::candidate_commit(
243 commit_ref.clone(),
244 &stage.commit.value.to_bytes(),
245 stage.commit.prepared.stored_bytes(),
246 )
247 .map_err(|error| DbError::context("prepared candidate commit", error))?;
248 persist_exact_remote_object_on(&tx, self.store_dir, &commit_remote, "candidate commit")?;
249 let expected_partition_count = usize::from(partitions.store.is_some())
250 .checked_add(partitions.circles.len())
251 .ok_or_else(|| DbError::Message("audience partition count overflow".to_string()))?;
252 if stage.audiences.packages.len() != expected_partition_count {
253 return Err(DbError::Message(
254 "prepared audience packages do not cover every write partition".to_string(),
255 ));
256 }
257 let mut indexed = std::collections::BTreeSet::new();
258 for package in &stage.audiences.packages {
259 let value = package.package();
260 if value.store_root_hash() != stage.root.store_root_hash
261 || value.write_id() != &stage.write_id
262 || value.commit_coord() != &commit_ref.coord
263 || value.candidate_family() != stage.commit.value.candidate_family()
264 {
265 return Err(DbError::Message(
266 "prepared audience package differs from its exact Store commit".to_string(),
267 ));
268 }
269 match value.audience() {
270 coven_protocol::audience_package::PackageAudience::Store => {
271 let partition = partitions.store.as_ref().ok_or_else(|| {
272 DbError::Message(
273 "prepared Store package has no Store partition".to_string(),
274 )
275 })?;
276 if value.changeset() != partition.changeset {
277 return Err(DbError::Message(
278 "prepared Store package changeset differs from its partition"
279 .to_string(),
280 ));
281 }
282 stage
283 .commit
284 .value
285 .verify_store_package(package.semantic_bytes())
286 .map_err(DbError::from)?;
287 }
288 coven_protocol::audience_package::PackageAudience::Circle { circle_id, .. } => {
289 let partition = partitions
290 .circles
291 .iter()
292 .find(|partition| {
293 partition.audience
294 == coven_protocol::circle::Audience::Circle(*circle_id)
295 })
296 .ok_or_else(|| {
297 DbError::Message(format!(
298 "prepared Circle package {circle_id} has no partition"
299 ))
300 })?;
301 if value.changeset() != partition.changeset {
302 return Err(DbError::Message(format!(
303 "prepared Circle package {circle_id} changeset differs from its partition"
304 )));
305 }
306 stage
307 .commit
308 .value
309 .verify_circle_package(*circle_id, package.semantic_bytes())
310 .map_err(DbError::from)?;
311 }
312 }
313 indexed.insert(package.remote_object_id());
314 }
315 indexed.extend(
316 stage
317 .audiences
318 .blobs
319 .iter()
320 .map(crate::PreparedAudienceBlob::remote_object_id),
321 );
322 debug_assert_eq!(indexed, object_ids);
323 super::prepared_remote_objects::persist_prepared_audience_objects_on(
324 &tx,
325 self.store_dir,
326 &stage.write_id,
327 &stage.audiences.packages,
328 &stage.audiences.blobs,
329 )?;
330 let prepared = PreparedStoreWriteState {
331 commit: DurablePreparedProtocolObject::new(
332 stage.commit.value.to_bytes(),
333 stage.commit.prepared,
334 ),
335 history_evidence: stage.history_evidence,
336 local_cleanup: stage.local_cleanup,
337 completion: stage.completion,
338 };
339 let prepared = serde_json::to_string(&prepared)
340 .map_err(|error| DbError::context("serialize prepared Store write", error))?;
341 let status = serde_json::to_string(&WriteStatus::Publishing)
342 .map_err(|error| DbError::context("serialize write status", error))?;
343 let updated = tx
344 .execute(
345 "UPDATE store_writes SET prepared = ?2, status = ?3
346 WHERE write_id = ?1 AND prepared IS NULL AND status = '\"pending\"'",
347 rusqlite::params![stage.write_id.as_str(), prepared, status],
348 )
349 .map_err(DbError::from)?;
350 if updated != 1 {
351 return Err(DbError::Message(format!(
352 "write {} lost pending preparation ownership",
353 stage.write_id
354 )));
355 }
356 if let Some(reserved) = reserved {
357 for retired in reserved.retired_candidates() {
358 super::candidate_records::begin_candidate_nonactivation_targets_on(
359 &tx,
360 &retired.candidate()?,
361 &retired.objects()?,
362 &retired.nonactivation,
363 )?;
364 }
365 }
366 tx.commit().map_err(DbError::from)?;
367 Ok(())
368 }
369
370 #[cfg(any(test, feature = "test-utils"))]
371 fn enqueue_store_changeset_for_test(
372 &mut self,
373 write_id: coven_protocol::write::WriteId,
374 changeset: Vec<u8>,
375 ) -> Result<(), DbError> {
376 let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
377 let changeset = self.blob_decls.complete_blob_changeset(&tx, &changeset)?;
378 let base = StoreWriteBase {
379 dependencies: crate::store::materialized_commit_index::materialized_frontier_on(
380 &tx, None,
381 )?,
382 };
383 let partitions = vec![crate::AudiencePartition {
384 audience: coven_protocol::circle::Audience::Store,
385 control: None,
386 changeset: changeset.clone(),
387 }];
388 let blob_facts = super::host_write_capture::capture_partition_blob_facts_on(
389 &tx,
390 &partitions,
391 self.blob_decls,
392 )?;
393 let changeset_hash =
394 crate::payload_store::write_payload_blocking(&tx, self.store_dir, &changeset)?;
395 crate::store::store_session::StoreTransaction::new(&tx, self.store_dir)
396 .insert_store_write(&write_id, &partitions, changeset_hash, &base, &blob_facts)?;
397 tx.commit().map_err(DbError::from)
398 }
399}
400
401impl StoreDatabase {
402 pub async fn table_schema_for_apply(&self) -> Result<crate::TableSchema, DbError> {
403 self.call_store(|session| session.table_schema_for_apply())
404 .await
405 }
406
407 pub async fn prepare_store_write_commit(
408 &self,
409 stage: StoreWritePreparation,
410 ) -> Result<(), DbError> {
411 let write_id = stage.write_id.clone();
412 self.call_store(move |session| session.prepare_store_write_commit(stage))
413 .await?;
414 self.notify_write_status(write_id, WriteStatus::Publishing);
415 Ok(())
416 }
417
418 #[cfg(any(test, feature = "test-utils"))]
419 pub async fn enqueue_store_changeset_for_test(
420 &self,
421 changeset: Vec<u8>,
422 ) -> Result<(), DbError> {
423 let write_id = self.new_store_write_id();
424 self.call_store(move |session| {
425 session.enqueue_store_changeset_for_test(write_id, changeset)
426 })
427 .await
428 }
429}