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 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 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 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}