Skip to main content

coven_database/store/store_session/
snapshot_retirement.rs

1use super::verified_store_authority::VerifiedRegistrationLookup;
2use super::*;
3use coven_protocol::objects::{ExactObjectRef, ObjectSlot};
4use coven_protocol::remote_object::{remote_object_id, SnapshotObjectOwner};
5use coven_protocol::store_commit::{AcceptedStoreSnapshotRef, SnapshotMeta, StoreRootRef};
6use std::collections::BTreeSet;
7
8impl StoreTransaction<'_, '_> {
9    /// The accepted successor owns every unfinished physical deletion. Release
10    /// local leases in the same transaction as adoption, before any remote I/O.
11    pub(super) fn retire_snapshot_artifact_ownership(
12        self,
13        lookup: &mut dyn VerifiedRegistrationLookup,
14        root: &StoreRootRef,
15        accepted: &AcceptedStoreSnapshotRef,
16        metadata: &SnapshotMeta,
17    ) -> Result<BTreeSet<ObjectSlot>, DbError> {
18        if metadata.store_root_hash != root.store_root_hash
19            || accepted.publication.store_root_hash != root.store_root_hash
20            || metadata.snapshot_hash() != accepted.snapshot.snapshot_hash
21            || metadata.publication_predecessor.next_position()? != accepted.publication.position
22        {
23            return Err(DbError::Message(
24                "snapshot retirement differs from its accepted boundary".into(),
25            ));
26        }
27        metadata
28            .history_summary
29            .reclaim
30            .validate_before(&accepted.publication)?;
31        let protected = metadata
32            .history_summary
33            .pending_device_join_snapshot_slots();
34        let mut superseded = metadata
35            .history_summary
36            .reclaim
37            .snapshots
38            .values()
39            .map(|snapshot| snapshot.accepted.snapshot.object.slot().clone())
40            .filter(|slot| !protected.contains(slot))
41            .collect::<BTreeSet<_>>();
42        for old in StoreRecords::new(self.transaction, self.store_dir)
43            .published_store_snapshots(root, lookup)?
44        {
45            let position = old.meta.publication_predecessor.next_position()?;
46            if position >= accepted.publication.position
47                || protected.contains(old.reference.object.slot())
48            {
49                continue;
50            }
51            superseded.insert(old.reference.object.slot().clone());
52            let owner = SnapshotObjectOwner::Store {
53                metadata_slot: old.reference.object.slot().clone(),
54            };
55            for (object, image) in [
56                (&old.meta.image.object, true),
57                (&old.meta.membership_rollup.object, false),
58            ] {
59                let object_id = remote_object_id(object);
60                let exists: bool = self.transaction.query_row(
61                    "SELECT EXISTS(SELECT 1 FROM remote_objects WHERE object_id = ?1)",
62                    [object_id.to_string()],
63                    |row| row.get(0),
64                )?;
65                if !exists {
66                    continue;
67                }
68                let remote = crate::remote_object_records::load_remote_object_on(
69                    self.transaction,
70                    object_id,
71                )?;
72                if image {
73                    remote.validate_reclaimable_snapshot_image(&old.meta.image, &owner)
74                } else {
75                    remote
76                        .validate_reclaimable_membership_rollup(&old.meta.membership_rollup, &owner)
77                }
78                .map_err(|error| {
79                    DbError::context("retire superseded snapshot artifact lease", error)
80                })?;
81                if !crate::remote_object_records::delete_remote_object_on(
82                    self.transaction,
83                    object_id,
84                )? {
85                    return Err(DbError::Message(
86                        "snapshot artifact lease disappeared during retirement".into(),
87                    ));
88                }
89            }
90            self.transaction.execute(
91                "DELETE FROM published_store_snapshot WHERE publication_position = ?1 AND snapshot_ref = ?2",
92                rusqlite::params![
93                    i64::try_from(position.get()).map_err(|error| DbError::context("retired snapshot position", error))?,
94                    serde_json::to_string(&old.reference).map_err(|error| DbError::context("retired snapshot reference", error))?,
95                ],
96            )?;
97        }
98        Ok(superseded)
99    }
100
101    pub(super) fn retire_store_blob_snapshot_ownership(
102        self,
103        metadata_slot: &ObjectSlot,
104        superseded: &BTreeSet<ObjectSlot>,
105    ) -> Result<(), DbError> {
106        let mut statement = self
107            .transaction
108            .prepare("SELECT remote_object_id FROM blob_locators ORDER BY remote_object_id")?;
109        let blobs = statement
110            .query_map([], |row| row.get::<_, String>(0))?
111            .collect::<Result<Vec<_>, _>>()?;
112        drop(statement);
113        for blob in blobs {
114            let object_id = blob
115                .parse()
116                .map_err(|error| DbError::context("snapshot blob object id", error))?;
117            let mut remote =
118                crate::remote_object_records::load_remote_object_on(self.transaction, object_id)?;
119            if !remote.is_activated_stored_blob()
120                || !remote.snapshot_owners().any(|owner| {
121                    matches!(owner, SnapshotObjectOwner::Store { metadata_slot }
122                    if superseded.contains(metadata_slot))
123                })
124            {
125                continue;
126            }
127            remote
128                .retire_superseded_store_snapshot_ownership(metadata_slot, superseded)
129                .map_err(|error| {
130                    DbError::context("retire accepted snapshot blob ownership", error)
131                })?;
132            crate::remote_object_records::update_remote_object_on(
133                self.transaction,
134                object_id,
135                &remote,
136            )?;
137        }
138        Ok(())
139    }
140}
141
142impl StoreDatabase {
143    pub async fn verify_snapshot_artifacts_released(
144        &self,
145        objects: Vec<ExactObjectRef>,
146    ) -> Result<(), DbError> {
147        self.call_store(move |session| {
148            for object in objects {
149                let id = remote_object_id(&object);
150                let exists: bool = session.conn.query_row(
151                    "SELECT EXISTS(SELECT 1 FROM remote_objects WHERE object_id = ?1)",
152                    [id.to_string()],
153                    |row| row.get(0),
154                )?;
155                if exists {
156                    crate::remote_object_records::load_remote_object_on(session.conn, id)?;
157                    return Err(DbError::Message(
158                        "accepted snapshot artifact retains a live local owner".into(),
159                    ));
160                }
161            }
162            Ok(())
163        })
164        .await
165    }
166}