Skip to main content

mkit_server/pipeline/
purge.rs

1use super::{Pipeline, meta_error};
2use crate::pipeline::HookSet;
3#[cfg(test)]
4use crate::purge::plan_enqueue;
5use crate::purge::{Request, Trigger};
6#[cfg(test)]
7use crate::store::keys;
8use crate::store::{Batch, MultipartBlobStore, NamespaceStore, Partition};
9use crate::{RepoId, ServerError};
10
11impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
12    /// Plan an audited safety purge in the same batch as a serving stop.
13    /// `operation_id` is stable across retries of the automatic action.
14    /// Callers merge this batch into the authoritative state mutation.
15    pub async fn plan_repository_purge(
16        &self,
17        partition: &Partition,
18        repo: &RepoId,
19        trigger: Trigger,
20        operation_id: &str,
21        now_ms: u64,
22    ) -> Result<Batch, ServerError> {
23        crate::purge::automatic::plan_repository(
24            self.cfg.purge.as_ref(),
25            &self.meta,
26            partition,
27            repo,
28            trigger,
29            operation_id,
30            now_ms,
31        )
32        .await
33        .map_err(meta_error)
34    }
35    /// Immediately invalidate local entries after committing a serving stop and
36    /// its durable purge work. Timer 11 retries any failed cache deletion.
37    pub async fn invalidate_local_cache(&self, repo: &RepoId) {
38        #[cfg(feature = "http-objects")]
39        if let Some(seams) = &self.http_seams {
40            seams.reachability.invalidate(repo);
41        }
42        if let Some(config) = &self.cfg.purge
43            && let Some(local) = &config.local
44        {
45            let request = Request {
46                purge_id: "local".into(),
47                audience: config.audience.clone(),
48                repository: format!("{}/{}", repo.namespace.as_str(), repo.name.as_str()),
49                namespace: String::new(),
50                trigger: Trigger::Suspension,
51                url_paths: Vec::new(),
52                object_ids: Vec::new(),
53                refs: Vec::new(),
54            };
55            // The state change, immutable purge work and timer already committed.
56            // Failed or incomplete immediate deletes remain pending in kind 11.
57            let _ = local
58                .invalidate(&request, 0, &crate::purge::SliceBudget::new(64))
59                .await;
60        }
61    }
62}
63
64#[cfg(all(test, feature = "memory"))]
65mod tests {
66    use super::*;
67    use crate::purge::{LocalInvalidation, PurgeConfig, SliceBudget};
68    use crate::{
69        Addressing, ManualClock, MemoryBlobStore, MemoryKv, NamespaceKey, RepoName, StoreError,
70    };
71    use std::sync::{Arc, Mutex};
72
73    #[tokio::test]
74    async fn automatic_purge_identity_is_source_partition_bound_and_retry_stable() {
75        let store = Arc::new(MemoryKv::default());
76        let repo = RepoId {
77            namespace: NamespaceKey::deployment_default(),
78            name: RepoName::new("repo").unwrap(),
79        };
80        let coordinator = Partition::Coordinator(repo.namespace.clone());
81        let index = Partition::RepoIndex {
82            ns: repo.namespace.clone(),
83            repo: repo.name.clone(),
84            prefix: 1,
85        };
86        let mut config = crate::pipeline::PipelineConfig::new(
87            Addressing::Single { repo: repo.clone() },
88            crate::pipeline::AuthMode::Open,
89            crate::upload::UploadLimits {
90                max_total_bytes: 1024,
91                max_chunks: 32,
92            },
93        );
94        config.purge = Some(
95            PurgeConfig::new("https://server.example".into(), true, true).with_audit(Arc::new(
96                crate::admin::SystemAudit::new(store.clone(), coordinator.clone()),
97            )),
98        );
99        let pipeline = Pipeline::new(
100            MemoryBlobStore::default(),
101            store,
102            crate::pipeline::Hooks::new(),
103            config,
104            Arc::new(ManualClock::new(10)),
105            Arc::new(crate::NoopMetrics),
106        )
107        .unwrap();
108        let mut requests = Vec::new();
109        for partition in [&coordinator, &index, &coordinator] {
110            let batch = pipeline
111                .plan_repository_purge(partition, &repo, Trigger::Suspension, "same-op", 10)
112                .await
113                .unwrap();
114            requests.push(
115                batch
116                    .writes
117                    .iter()
118                    .find_map(|write| match write {
119                        crate::Write::Put(key, value) if key.as_bytes().starts_with(b"cp\0") => {
120                            Some(
121                                serde_json::from_slice::<Request>(value.as_bytes())
122                                    .expect("planned purge request"),
123                            )
124                        }
125                        _ => None,
126                    })
127                    .unwrap(),
128            );
129        }
130        assert_ne!(requests[0].purge_id, requests[1].purge_id);
131        assert_eq!(requests[0].purge_id, requests[2].purge_id);
132    }
133    struct FailingLocal {
134        store: Arc<MemoryKv>,
135        partition: Partition,
136        seen: Arc<Mutex<Vec<Request>>>,
137    }
138    impl LocalInvalidation for FailingLocal {
139        fn invalidate<'a>(
140            &'a self,
141            request: &'a Request,
142            _: u32,
143            _: &'a SliceBudget,
144        ) -> crate::BoxFuture<'a, Result<Option<u32>, StoreError>> {
145            Box::pin(async move {
146                assert!(
147                    crate::purge::read_request(&self.store, &self.partition, "committed")
148                        .await?
149                        .is_some()
150                );
151                self.seen
152                    .lock()
153                    .expect("local invalidation recording lock")
154                    .push(request.clone());
155                Err(StoreError::unavailable("local cache offline"))
156            })
157        }
158    }
159    #[tokio::test]
160    async fn immediate_local_failure_keeps_committed_intent_and_timer_for_retry() {
161        let store = Arc::new(MemoryKv::default());
162        let repo = RepoId {
163            namespace: NamespaceKey::deployment_default(),
164            name: RepoName::new("repo").unwrap(),
165        };
166        let partition = Partition::Namespace(repo.namespace.clone());
167        let request = Request {
168            purge_id: "committed".into(),
169            audience: "https://server.example".into(),
170            repository: "root/repo".into(),
171            namespace: String::new(),
172            trigger: Trigger::Suspension,
173            url_paths: Vec::new(),
174            object_ids: Vec::new(),
175            refs: Vec::new(),
176        };
177        store
178            .apply(&partition, plan_enqueue(&request, 10, None, None).unwrap())
179            .await
180            .unwrap();
181        let seen = Arc::new(Mutex::new(Vec::new()));
182        let mut config = crate::pipeline::PipelineConfig::new(
183            Addressing::Single { repo: repo.clone() },
184            crate::pipeline::AuthMode::Open,
185            crate::upload::UploadLimits {
186                max_total_bytes: 1024,
187                max_chunks: 32,
188            },
189        );
190        config.purge = Some(
191            PurgeConfig::new(request.audience.clone(), true, true).with_local(Arc::new(
192                FailingLocal {
193                    store: store.clone(),
194                    partition: partition.clone(),
195                    seen: seen.clone(),
196                },
197            )),
198        );
199        let pipeline = Pipeline::new(
200            MemoryBlobStore::default(),
201            store.clone(),
202            crate::pipeline::Hooks::new(),
203            config,
204            Arc::new(ManualClock::new(10)),
205            Arc::new(crate::NoopMetrics),
206        )
207        .unwrap();
208        pipeline.invalidate_local_cache(&repo).await;
209        assert_eq!(seen.lock().unwrap()[0].repository, "root/repo");
210        assert!(
211            crate::purge::read_request(&store, &partition, "committed")
212                .await
213                .unwrap()
214                .is_some()
215        );
216        assert!(
217            store
218                .get(&partition, &keys::timer(10, 11, b"committed"))
219                .await
220                .unwrap()
221                .is_some()
222        );
223    }
224}