Skip to main content

coven_database/store/store_session/
preparation.rs

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(&registration_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            &registration_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                &registration.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}