This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-maintenance / src / scheduler.rs
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}