This repository has no description
17 kB
515 lines
1use std::collections::HashSet;
2use std::sync::Arc;
3use std::time::{Duration, SystemTime};
4
5use knot_git::Layout;
6use knot_lfs::DiskStore;
7use knot_runtime::Clock;
8use knot_types::{RepoDid, UnixSeconds};
9use tokio::sync::{mpsc, watch};
10
11use crate::{LfsGrace, Options, RepackStatus, Report, SweepInterval, run_repo};
12
13pub trait RepoSource: Send + Sync + 'static {
14 fn repos(&self) -> Vec<RepoDid>;
15 fn ready_repos(&self) -> Option<Vec<RepoDid>>;
16}
17
18// "Grace" means how old an unreferenced object can get before we reap it,
19// and the "interval" is how often we look.
20struct LfsGc {
21 store: Arc<DiskStore>,
22 grace: LfsGrace,
23 interval: SweepInterval,
24}
25
26const TRIGGER_CAPACITY: usize = 1024;
27
28knot_types::scalar_newtype! {
29 pub struct PushBytes(u64) => ordered;
30}
31
32#[derive(Clone)]
33pub struct MaintenanceHandle {
34 trigger: Option<mpsc::Sender<RepoDid>>,
35 large_push: PushBytes,
36}
37
38impl MaintenanceHandle {
39 pub fn disabled() -> Self {
40 Self {
41 trigger: None,
42 large_push: PushBytes::new(u64::MAX),
43 }
44 }
45
46 pub fn note_push(&self, repo: &RepoDid, pack_bytes: PushBytes) {
47 if pack_bytes >= self.large_push
48 && let Some(trigger) = &self.trigger
49 {
50 let _ = trigger.try_send(repo.clone());
51 }
52 }
53}
54
55pub struct Scheduler<C> {
56 layout: Layout,
57 source: Arc<dyn RepoSource>,
58 clock: C,
59 options: Options,
60 interval: Duration,
61 triggers: mpsc::Receiver<RepoDid>,
62 lfs: Option<LfsGc>,
63}
64
65impl<C: Clock> Scheduler<C> {
66 pub fn new(
67 layout: Layout,
68 source: Arc<dyn RepoSource>,
69 clock: C,
70 options: Options,
71 interval: Duration,
72 large_push: PushBytes,
73 ) -> (Self, MaintenanceHandle) {
74 let (trigger, triggers) = mpsc::channel(TRIGGER_CAPACITY);
75 let handle = MaintenanceHandle {
76 trigger: Some(trigger),
77 large_push,
78 };
79 let scheduler = Self {
80 layout,
81 source,
82 clock,
83 options,
84 interval,
85 triggers,
86 lfs: None,
87 };
88 (scheduler, handle)
89 }
90
91 pub fn with_lfs_gc(
92 mut self,
93 store: Arc<DiskStore>,
94 grace: LfsGrace,
95 interval: SweepInterval,
96 ) -> Self {
97 self.lfs = Some(LfsGc {
98 store,
99 grace,
100 interval,
101 });
102 self
103 }
104
105 fn now_seconds(&self) -> UnixSeconds {
106 UnixSeconds::new((self.clock.now_unix_micros().get() / 1_000_000) as i64)
107 }
108
109 fn now(&self) -> SystemTime {
110 SystemTime::UNIX_EPOCH + Duration::from_micros(self.clock.now_unix_micros().get())
111 }
112
113 async fn gc_repo(&self, repo: &RepoDid) -> knot_lfs::GcReport {
114 let Some(lfs) = &self.lfs else {
115 return knot_lfs::GcReport::default();
116 };
117 let layout = self.layout.clone();
118 let store = Arc::clone(&lfs.store);
119 let grace = lfs.grace.get();
120 let now = self.now();
121 let target = repo.clone();
122 let started = std::time::Instant::now();
123 let outcome = tokio::task::spawn_blocking(move || {
124 layout
125 .open(&target)
126 .map_err(knot_lfs::GcError::from)
127 .and_then(|opened| knot_lfs::collect_repo(&store, &opened, &target, grace, now))
128 })
129 .await;
130 match outcome {
131 Ok(Ok(report)) => {
132 if report.swept > 0 {
133 tracing::info!(
134 repo = %repo,
135 scanned = report.scanned,
136 marked = report.marked,
137 swept = report.swept,
138 bytes = report.bytes.get(),
139 duration_ms = started.elapsed().as_millis() as u64,
140 "lfs gc reclaimed objects"
141 );
142 }
143 report
144 }
145 Ok(Err(error)) => {
146 tracing::warn!(repo = %repo, %error, "lfs gc skipped, projection uncertain");
147 knot_lfs::GcReport::default()
148 }
149 Err(join) => {
150 tracing::error!(repo = %repo, %join, "lfs gc task panicked");
151 knot_lfs::GcReport::default()
152 }
153 }
154 }
155
156 async fn sweep_orphans(&self) {
157 let Some(lfs) = &self.lfs else {
158 return;
159 };
160 let Some(hosted) = self.source.ready_repos() else {
161 return;
162 };
163 let store = Arc::clone(&lfs.store);
164 let grace = lfs.grace.get();
165 let now = self.now();
166 let hosted: HashSet<RepoDid> = hosted.into_iter().collect();
167 let outcome =
168 tokio::task::spawn_blocking(move || store.sweep_orphans(&hosted, grace, now)).await;
169 match outcome {
170 Ok(Ok(sweep)) if sweep.prefixes > 0 => tracing::info!(
171 prefixes = sweep.prefixes,
172 objects = sweep.objects,
173 bytes = sweep.bytes.get(),
174 "lfs gc reclaimed orphan prefixes"
175 ),
176 Ok(Ok(_)) => {}
177 Ok(Err(error)) => tracing::warn!(%error, "lfs orphan sweep failed"),
178 Err(join) => tracing::error!(%join, "lfs orphan sweep task panicked"),
179 }
180 }
181
182 pub async fn run(mut self, mut shutdown: watch::Receiver<bool>) {
183 let mut tick = tokio::time::interval(self.interval);
184 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
185 tick.tick().await;
186 let mut lfs_tick = self.lfs.as_ref().map(|lfs| {
187 let mut lfs_tick = tokio::time::interval(lfs.interval.get());
188 lfs_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
189 lfs_tick
190 });
191 if let Some(lfs_tick) = &mut lfs_tick {
192 lfs_tick.tick().await;
193 }
194 let probe = shutdown.clone();
195 loop {
196 tokio::select! {
197 biased;
198 _ = shutdown.changed() => break,
199 _ = tick.tick() => self.sweep_all(&probe).await,
200 _ = tick_lfs(&mut lfs_tick) => self.lfs_sweep_all(&probe).await,
201 Some(repo) = self.triggers.recv() => self.maintain_batch(repo).await,
202 }
203 }
204 }
205
206 async fn sweep_all(&self, shutdown: &watch::Receiver<bool>) {
207 for repo in self.source.repos() {
208 if *shutdown.borrow() {
209 return;
210 }
211 self.maintain_one(repo).await;
212 }
213 }
214
215 async fn lfs_sweep_all(&self, shutdown: &watch::Receiver<bool>) {
216 if self.lfs.is_none() {
217 return;
218 }
219 let started = std::time::Instant::now();
220 let mut totals = knot_lfs::GcReport::default();
221 let mut repos = 0usize;
222 for repo in self.source.repos() {
223 if *shutdown.borrow() {
224 return;
225 }
226 let report = self.gc_repo(&repo).await;
227 repos += 1;
228 totals.scanned += report.scanned;
229 totals.marked += report.marked;
230 totals.swept += report.swept;
231 totals.bytes = totals.bytes.saturating_add(report.bytes);
232 }
233 if !*shutdown.borrow() {
234 self.sweep_orphans().await;
235 }
236 tracing::info!(
237 repos,
238 scanned = totals.scanned,
239 marked = totals.marked,
240 swept = totals.swept,
241 bytes = totals.bytes.get(),
242 duration_ms = started.elapsed().as_millis() as u64,
243 "lfs gc pass finished"
244 );
245 }
246
247 async fn maintain_batch(&mut self, first: RepoDid) {
248 let pending: HashSet<RepoDid> = std::iter::once(first)
249 .chain(std::iter::from_fn(|| self.triggers.try_recv().ok()))
250 .collect();
251 for repo in pending {
252 self.maintain_one(repo.clone()).await;
253 self.gc_repo(&repo).await;
254 }
255 }
256
257 async fn maintain_one(&self, repo: RepoDid) {
258 let layout = self.layout.clone();
259 let options = self.options;
260 let now_seconds = self.now_seconds();
261 let target = repo.clone();
262 let outcome = tokio::task::spawn_blocking(move || {
263 layout
264 .open(&target)
265 .map_err(crate::MaintError::from)
266 .and_then(|opened| run_repo(&opened, now_seconds, &options))
267 })
268 .await;
269 match outcome {
270 Ok(Ok(report)) => report_skips(&repo, &report),
271 Ok(Err(error)) => tracing::error!(repo = %repo, %error, "maintenance run failed"),
272 Err(join) => tracing::error!(repo = %repo, %join, "maintenance task panicked"),
273 }
274 }
275}
276
277async fn tick_lfs(tick: &mut Option<tokio::time::Interval>) {
278 match tick {
279 Some(tick) => {
280 tick.tick().await;
281 }
282 None => std::future::pending::<()>().await,
283 }
284}
285
286fn report_skips(repo: &RepoDid, report: &Report) {
287 match report.repack.status {
288 RepackStatus::SkippedTooLarge => {
289 tracing::warn!(
290 repo = %repo,
291 reason = "reachable set exceeds repack_max_objects",
292 "skipped repack"
293 )
294 }
295 RepackStatus::ClosureFailed => {
296 tracing::warn!(
297 repo = %repo,
298 reason = "couldn't compute reachable set",
299 "skipped repack and prune"
300 )
301 }
302 _ => {}
303 }
304}
305
306#[cfg(test)]
307mod tests {
308 use std::sync::Arc;
309
310 use knot_git::{Layout, RefUpdate};
311 use knot_runtime::SystemClock;
312 use knot_types::{BranchName, RefName, RepoDid};
313
314 use super::{MaintenanceHandle, PushBytes, RepoSource, Scheduler};
315 use crate::test_support::{commit_on, empty_tree, options};
316
317 struct Fixed(Vec<RepoDid>);
318 impl RepoSource for Fixed {
319 fn repos(&self) -> Vec<RepoDid> {
320 self.0.clone()
321 }
322
323 fn ready_repos(&self) -> Option<Vec<RepoDid>> {
324 Some(self.0.clone())
325 }
326 }
327
328 fn seed_repo(layout: &Layout, did: &RepoDid) {
329 let repo = layout.create(did).unwrap();
330 let tip = commit_on(&repo, empty_tree(repo.object_format()), Vec::new(), "a");
331 repo.update_ref(&RefUpdate::Create {
332 name: RefName::new("refs/heads/main").unwrap(),
333 new: tip,
334 })
335 .unwrap();
336 }
337
338 #[test]
339 fn note_push_fires_only_past_the_threshold() {
340 let (sender, mut receiver) = tokio::sync::mpsc::channel(16);
341 let handle = MaintenanceHandle {
342 trigger: Some(sender),
343 large_push: PushBytes::new(1_000),
344 };
345 let did = RepoDid::new("did:plc:squid").unwrap();
346 handle.note_push(&did, PushBytes::new(999));
347 assert!(receiver.try_recv().is_err(), "small push is ignored");
348 handle.note_push(&did, PushBytes::new(1_000));
349 assert_eq!(receiver.try_recv().unwrap(), did, "large push triggers");
350 }
351
352 #[test]
353 fn disabled_handle_never_triggers() {
354 let did = RepoDid::new("did:plc:squid").unwrap();
355 MaintenanceHandle::disabled().note_push(&did, PushBytes::new(u64::MAX));
356 }
357
358 #[tokio::test]
359 async fn the_lfs_interval_collects_without_a_push_trigger() {
360 use std::time::{Duration, SystemTime, UNIX_EPOCH};
361
362 use knot_lfs::{DiskStore, LfsOid, LfsStore, LfsStorePath};
363 use knot_runtime::{ManualClock, UnixMicros};
364 use sha2::{Digest, Sha256};
365
366 let scan = tempfile::tempdir().unwrap();
367 let lfs_dir = tempfile::tempdir().unwrap();
368 let layout = Layout::new(scan.path()).with_default_branch(BranchName::new("main").unwrap());
369 let did = RepoDid::new("did:plc:limpet").unwrap();
370 seed_repo(&layout, &did);
371
372 let store = Arc::new(DiskStore::open(LfsStorePath::new(lfs_dir.path())).unwrap());
373 let body: &[u8] = b"unreferenced media reclaimed on the interval alone";
374 let oid = LfsOid::from_digest(Sha256::digest(body).into());
375 let size = knot_lfs::ClaimedSize::new(body.len() as u64);
376 store.put(&did, &oid, size, &mut &body[..]).unwrap();
377
378 let real_micros = SystemTime::now()
379 .duration_since(UNIX_EPOCH)
380 .unwrap()
381 .as_micros() as u64;
382 let future = ManualClock::new(UnixMicros::new(real_micros + 5 * 86_400 * 1_000_000));
383
384 let (scheduler, _handle) = Scheduler::new(
385 layout.clone(),
386 Arc::new(Fixed(vec![did.clone()])),
387 future,
388 options(),
389 std::time::Duration::from_secs(3_600),
390 PushBytes::new(1_000),
391 );
392 let scheduler = scheduler.with_lfs_gc(
393 Arc::clone(&store),
394 crate::lfs_grace(
395 crate::GcGrace::from_secs(86_400),
396 crate::ReflogRetention::from_secs(90 * 86_400),
397 ),
398 crate::SweepInterval::new(Duration::from_millis(40)),
399 );
400 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
401 let task = tokio::spawn(scheduler.run(shutdown_rx));
402
403 let mut waited = 0;
404 while store.probe(&did, &oid).unwrap().is_some() && waited < 200 {
405 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
406 waited += 1;
407 }
408 assert_eq!(
409 store.probe(&did, &oid).unwrap(),
410 None,
411 "the lfs interval alone reclaimed the unreferenced, past-grace object"
412 );
413
414 shutdown_tx.send(true).unwrap();
415 task.await.unwrap();
416 }
417
418 #[tokio::test]
419 async fn a_triggered_push_collects_an_unreferenced_lfs_object() {
420 use std::time::{Duration, SystemTime, UNIX_EPOCH};
421
422 use knot_lfs::{DiskStore, LfsOid, LfsStore, LfsStorePath};
423 use knot_runtime::{ManualClock, UnixMicros};
424 use sha2::{Digest, Sha256};
425
426 let scan = tempfile::tempdir().unwrap();
427 let lfs_dir = tempfile::tempdir().unwrap();
428 let layout = Layout::new(scan.path()).with_default_branch(BranchName::new("main").unwrap());
429 let did = RepoDid::new("did:plc:limpet").unwrap();
430 seed_repo(&layout, &did);
431
432 let store = Arc::new(DiskStore::open(LfsStorePath::new(lfs_dir.path())).unwrap());
433 let body: &[u8] = b"unreferenced media the sweep should reclaim";
434 let oid = LfsOid::from_digest(Sha256::digest(body).into());
435 let size = knot_lfs::ClaimedSize::new(body.len() as u64);
436 store.put(&did, &oid, size, &mut &body[..]).unwrap();
437
438 let real_micros = SystemTime::now()
439 .duration_since(UNIX_EPOCH)
440 .unwrap()
441 .as_micros() as u64;
442 let future = ManualClock::new(UnixMicros::new(real_micros + 5 * 86_400 * 1_000_000));
443
444 let (scheduler, handle) = Scheduler::new(
445 layout.clone(),
446 Arc::new(Fixed(vec![did.clone()])),
447 future,
448 options(),
449 std::time::Duration::from_secs(3_600),
450 PushBytes::new(1_000),
451 );
452 let scheduler = scheduler.with_lfs_gc(
453 Arc::clone(&store),
454 crate::lfs_grace(
455 crate::GcGrace::from_secs(86_400),
456 crate::ReflogRetention::from_secs(90 * 86_400),
457 ),
458 crate::SweepInterval::new(Duration::from_secs(3_600)),
459 );
460 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
461 let task = tokio::spawn(scheduler.run(shutdown_rx));
462
463 handle.note_push(&did, PushBytes::new(10_000));
464
465 let mut waited = 0;
466 while store.probe(&did, &oid).unwrap().is_some() && waited < 200 {
467 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
468 waited += 1;
469 }
470 assert_eq!(
471 store.probe(&did, &oid).unwrap(),
472 None,
473 "the triggered gc reclaimed the unreferenced, past-grace object"
474 );
475
476 shutdown_tx.send(true).unwrap();
477 task.await.unwrap();
478 }
479
480 #[tokio::test]
481 async fn a_triggered_repo_is_maintained() {
482 let scan = tempfile::tempdir().unwrap();
483 let layout = Layout::new(scan.path()).with_default_branch(BranchName::new("main").unwrap());
484 let did = RepoDid::new("did:plc:limpet").unwrap();
485 seed_repo(&layout, &did);
486
487 let (scheduler, handle) = Scheduler::new(
488 layout.clone(),
489 Arc::new(Fixed(vec![did.clone()])),
490 SystemClock,
491 options(),
492 std::time::Duration::from_secs(3_600),
493 PushBytes::new(1_000),
494 );
495 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
496 let task = tokio::spawn(scheduler.run(shutdown_rx));
497
498 handle.note_push(&did, PushBytes::new(10_000));
499
500 let graph = layout
501 .open(&did)
502 .unwrap()
503 .objects_dir()
504 .join("info/commit-graph");
505 let mut waited = 0;
506 while !graph.exists() && waited < 200 {
507 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
508 waited += 1;
509 }
510 assert!(graph.exists(), "triggered repo got commit-graph");
511
512 shutdown_tx.send(true).unwrap();
513 task.await.unwrap();
514 }
515}