Skip to content

g1t/services/actions/src/cache.rs

533 lines23,592 bytesCodeBlameRaw
1//! The cache of `actions/cache`: which entries each repository has, kept
2//! here, while the API keeps their bytes in R2 (the ACTIONS_CACHE bucket).
3//!
4//! - Entries are scoped by ref, as on GitHub (`scopes`): an entry belongs
5//! to the ref whose run saved it, and a run restores from its own ref,
6//! then its pull request's base branch, then the default branch. A pull
7//! request from outside the repository saves under `untrusted:<ref>`,
8//! which no other ref reads, so it can never plant an entry the default
9//! branch restores. g1t's own `actions/cache` and the toolkit's protocols
10//! follow the same rule (`Actions::cache_scope`).
11//! - An entry's version is the hash of its paths and compression, which
12//! the toolkit's client and g1t's runner both send: the same key saved
13//! for other paths is another entry.
14//! - A key is written once in its scope and version. In each scope in
15//! turn, a restore finds its key exactly, else the newest entry whose key
16//! starts with one of its restore keys.
17//! - An entry is at most `CACHE_MAX_ENTRY_BYTES`. A repository's entries
18//! hold at most `CACHE_REPO_QUOTA_BYTES` together: saving past it evicts
19//! the entries restored longest ago.
20//! - An entry not restored for `CACHE_UNUSED_DAYS`, or saved more than
21//! `CACHE_MAX_AGE_DAYS` ago, is deleted by the hourly sweep, which also
22//! writes down what each workspace holds, its artifacts included (they
23//! are in the same bucket), and reports this month's storage to billing
24//! (`note_pending`, source `cache`) at R2's price.
25//!
26//! The API calls these with the job's token, or its runtime token for the
27//! toolkit's protocols (runtime.rs), which is checked here. An entry the
28//! toolkit saves has the toolkit's version, and is found only by the same
29//! version; g1t's own `actions/cache` saves and finds entries without one.
30
31use g1t_contracts::FailureCode;
32use g1t_contracts::Outcome;
33use g1t_contracts::actions::{
34 CACHE_MAX_AGE_DAYS, CACHE_MAX_ENTRY_BYTES, CACHE_MICROS_PER_GB_MONTH, CACHE_REPO_QUOTA_BYTES, CACHE_UNUSED_DAYS,
35 CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs, CacheUploadArgs,
36};
37use g1t_contracts::billing::NotePendingArgs;
38use g1t_contracts::new_id;
39use g1t_contracts::time::rfc3339;
40use g1t_kit::now_ms;
41use serde::Deserialize;
42use serde_json::Value;
43use worker::Result;
44
45use crate::{Actions, check, fail};
46
47const DAY_MS: u64 = 24 * 60 * 60 * 1000;
48/// An upload not finished after this long is given up.
49const PENDING_MS: u64 = 6 * 60 * 60 * 1000;
50/// A GB, as storage is billed.
51const GB: f64 = 1_000_000_000.0;
52
53#[derive(Debug, Deserialize)]
54struct EntryRow {
55 id: String,
56 key: String,
57 object: String,
58 size: f64,
59 created_at: String,
60}
61
62/// The longest key: GitHub's limit.
63const MAX_KEY_CHARS: usize = 512;
64
65/// Whether a key can be kept: 1 to 512 characters, no commas (GitHub's rule).
66pub(crate) fn valid_key(key: &str) -> bool {
67 !key.is_empty() && key.chars().count() <= MAX_KEY_CHARS && !key.contains(',')
68}
69
70/// Which of `entries` (id, size, newest use first) to evict so that the
71/// repository holds at most `quota` with `keep` (just saved) among them.
72/// The entry just saved is never evicted.
73pub(crate) fn to_evict(entries: &[(String, u64)], keep: &str, quota: u64) -> Vec<String> {
74 let mut total: u64 = entries.iter().map(|(_, size)| size).sum();
75 let mut out = Vec::new();
76 for (id, size) in entries.iter().rev() {
77 if total <= quota {
78 break;
79 }
80 if id == keep {
81 continue;
82 }
83 total = total.saturating_sub(*size);
84 out.push(id.clone());
85 }
86 out
87}
88
89/// A month's GB-months from the bytes held each day so far: each day's
90/// bytes over 30 days.
91pub(crate) fn gb_months(days: &[u64]) -> f64 {
92 days.iter().map(|bytes| *bytes as f64).sum::<f64>() / GB / 30.0
93}
94
95/// What `gb_months` of cache cost g1t, in millionths of a dollar.
96pub(crate) fn storage_cost(gb_months: f64) -> i64 {
97 (gb_months * CACHE_MICROS_PER_GB_MONTH as f64).ceil() as i64
98}
99
100/// Where a job's cache entries are found and saved: its repository, and in
101/// it the refs it restores from, in order, and the one it saves to.
102pub(crate) struct CacheScope {
103 pub repo_id: String,
104 pub restore: Vec<String>,
105 pub save: String,
106}
107
108/// The scopes a run's jobs restore from, in order, and the one they save
109/// to: its own ref, then its pull request's base branch, then the default
110/// branch. A run that is not trusted (a pull request from outside) saves to
111/// a scope of its own that no other ref reads.
112pub(crate) fn scopes(git_ref: &str, base_ref: Option<&str>, default_branch: &str, trusted: bool) -> (Vec<String>, String) {
113 let own = if trusted { git_ref.to_owned() } else { format!("untrusted:{git_ref}") };
114 let mut restore = vec![own.clone()];
115 let base = base_ref.filter(|base| !base.is_empty()).map(|base| format!("refs/heads/{}", base.trim_start_matches("refs/heads/")));
116 for scope in base.into_iter().chain(std::iter::once(format!("refs/heads/{default_branch}"))) {
117 if !restore.contains(&scope) {
118 restore.push(scope);
119 }
120 }
121 (restore, own)
122}
123
124impl Actions {
125 /// Where a job's entries are found and saved (`scopes`), for g1t's own
126 /// `actions/cache` and the toolkit's protocols alike.
127 pub(crate) async fn cache_scope(&self, job: &crate::plan::JobRow) -> Result<CacheScope> {
128 let (restore, save) = match self.run_row(&job.run_id).await? {
129 Some(run) => {
130 let info = run.info();
131 scopes(&run.git_ref, info.base_ref.as_deref(), &info.default_branch, run.trusted != 0)
132 }
133 None => (Vec::new(), format!("untrusted:{}", job.run_id)),
134 };
135 Ok(CacheScope { repo_id: job.repo_id.clone(), restore, save })
136 }
137
138 /// A job by its own token or its runtime token (runtime.rs).
139 async fn cache_job(&self, job: &str, token: &str) -> Result<Outcome<crate::plan::JobRow>> {
140 self.job_for_credential(job, token).await
141 }
142
143 /// `cache_lookup`.
144 pub async fn cache_lookup(&self, a: CacheLookupArgs) -> Result<Outcome<Option<CacheHit>>> {
145 let job = check!(self.cache_job(&a.job, &a.token).await?);
146 let scope = self.cache_scope(&job).await?;
147 let now = now_ms();
148 let fresh = rfc3339(now.saturating_sub(CACHE_MAX_AGE_DAYS * DAY_MS));
149 // Found only by the same version: runners from before it was sent
150 // send none, and find only entries saved without one.
151 let version = a.version.clone().unwrap_or_default();
152 let mut found = None;
153 'scopes: for ref_scope in &scope.restore {
154 found = self
155 .db
156 .prepare(
157 "SELECT id, key, object, size, created_at FROM cache_entries
158 WHERE repo_id = ? AND scope = ? AND key = ? AND version = ? AND status = 'ready' AND created_at > ?",
159 )
160 .bind(&[
161 scope.repo_id.as_str().into(),
162 ref_scope.as_str().into(),
163 a.key.as_str().into(),
164 version.as_str().into(),
165 fresh.as_str().into(),
166 ])?
167 .first::<EntryRow>(None)
168 .await?;
169 if found.is_some() {
170 break;
171 }
172 for prefix in a.restore.iter().map(|p| p.trim()).filter(|p| !p.is_empty()) {
173 found = self
174 .db
175 .prepare(
176 "SELECT id, key, object, size, created_at FROM cache_entries
177 WHERE repo_id = ?1 AND scope = ?5 AND version = ?4 AND status = 'ready' AND created_at > ?3
178 AND substr(key, 1, length(?2)) = ?2
179 ORDER BY created_at DESC LIMIT 1",
180 )
181 .bind(&[
182 scope.repo_id.as_str().into(),
183 prefix.into(),
184 fresh.as_str().into(),
185 version.as_str().into(),
186 ref_scope.as_str().into(),
187 ])?
188 .first::<EntryRow>(None)
189 .await?;
190 if found.is_some() {
191 break 'scopes;
192 }
193 }
194 }
195 let Some(entry) = found else { return Ok(Outcome::Ok(None)) };
196 self.db
197 .prepare("UPDATE cache_entries SET last_used_at = ? WHERE id = ?")
198 .bind(&[rfc3339(now).into(), entry.id.as_str().into()])?
199 .run()
200 .await?;
201 let blob = a.version.as_ref().and_then(|_| self.download_token("cache", &entry.id, &entry.object, crate::runtime::DOWNLOAD_SECONDS));
202 Ok(Outcome::Ok(Some(CacheHit { key: entry.key, object: entry.object, size: entry.size as u64, created_at: entry.created_at, blob })))
203 }
204
205 /// `cache_upload`.
206 pub async fn cache_upload(&self, a: CacheUploadArgs) -> Result<Outcome<CacheReservation>> {
207 let job = check!(self.cache_job(&a.job, &a.token).await?);
208 let scope = self.cache_scope(&job).await?;
209 #[derive(Deserialize)]
210 struct Pending {
211 id: String,
212 object: String,
213 number: f64,
214 upload: Option<String>,
215 }
216 let columns = "SELECT id, object, rowid AS number, upload FROM cache_entries";
217 let pending = match (a.number, a.key.as_deref()) {
218 (Some(number), _) => {
219 self.db
220 .prepare(format!("{columns} WHERE rowid = ? AND repo_id = ? AND status = 'pending'"))
221 .bind(&[(number as f64).into(), scope.repo_id.as_str().into()])?
222 .first::<Pending>(None)
223 .await?
224 }
225 (None, Some(key)) => {
226 self.db
227 .prepare(format!("{columns} WHERE repo_id = ? AND scope = ? AND key = ? AND version = ? AND status = 'pending'"))
228 .bind(&[scope.repo_id.as_str().into(), scope.save.as_str().into(), key.into(), a.version.clone().unwrap_or_default().into()])?
229 .first::<Pending>(None)
230 .await?
231 }
232 (None, None) => None,
233 };
234 let Some(pending) = pending else {
235 return Ok(fail(FailureCode::NotFound, "No upload of that entry is in progress."));
236 };
237 let blob = match &pending.upload {
238 Some(upload) => match self.upload_token("cache", &pending.id, &pending.object, upload) {
239 Some(blob) => Some(blob),
240 None => return Ok(fail(FailureCode::Invalid, "The toolkit's storage is not set up here: the actions service has no ACTIONS_KEY.")),
241 },
242 None => None,
243 };
244 Ok(Outcome::Ok(CacheReservation { id: pending.id, object: pending.object, number: pending.number as u64, upload: pending.upload, blob }))
245 }
246
247 /// `cache_reserve`.
248 pub async fn cache_reserve(&self, a: CacheReserveArgs) -> Result<Outcome<CacheReservation>> {
249 let job = check!(self.cache_job(&a.job, &a.token).await?);
250 let scope = self.cache_scope(&job).await?;
251 if !valid_key(&a.key) {
252 return Ok(fail(FailureCode::Invalid, format!("A cache key is 1 to {MAX_KEY_CHARS} characters, without commas.")));
253 }
254 if a.size > CACHE_MAX_ENTRY_BYTES {
255 return Ok(fail(
256 FailureCode::Invalid,
257 format!("It is {} MB; a cache entry is at most {} MB.", a.size / 1_048_576, CACHE_MAX_ENTRY_BYTES / 1_048_576),
258 ));
259 }
260 let now = now_ms();
261 let at = rfc3339(now);
262 let version = a.version.clone().unwrap_or_default();
263 // An upload left unfinished long ago no longer holds its key.
264 self.db
265 .prepare(
266 "UPDATE cache_entries SET status = 'expired'
267 WHERE repo_id = ? AND scope = ? AND key = ? AND version = ? AND status = 'pending' AND created_at < ?",
268 )
269 .bind(&[
270 scope.repo_id.as_str().into(),
271 scope.save.as_str().into(),
272 a.key.as_str().into(),
273 version.as_str().into(),
274 rfc3339(now.saturating_sub(PENDING_MS)).into(),
275 ])?
276 .run()
277 .await?;
278 let id = new_id("cache", now);
279 let object = format!("c/{}/{id}", scope.repo_id);
280 #[derive(Deserialize)]
281 struct Inserted {
282 number: f64,
283 }
284 // A key is written once in its scope and version, never into
285 // another ref's scope.
286 let inserted = self
287 .db
288 .prepare(
289 "INSERT INTO cache_entries (id, repo_id, namespace, key, object, size, status, created_at, last_used_at, version, scope)
290 VALUES (?1, ?2, ?3, ?4, ?5, ?6, 'pending', ?7, ?7, ?8, ?9)
291 ON CONFLICT (repo_id, scope, key, version) DO UPDATE SET
292 id = ?1, object = ?5, size = ?6, status = 'pending', created_at = ?7, last_used_at = ?7, version = ?8, upload = NULL
293 WHERE cache_entries.status = 'expired'
294 RETURNING rowid AS number",
295 )
296 .bind(&[
297 id.as_str().into(),
298 scope.repo_id.as_str().into(),
299 job.namespace.as_str().into(),
300 a.key.as_str().into(),
301 object.as_str().into(),
302 (a.size as f64).into(),
303 at.as_str().into(),
304 version.as_str().into(),
305 scope.save.as_str().into(),
306 ])?
307 .first::<Inserted>(None)
308 .await?;
309 let Some(inserted) = inserted else {
310 return Ok(fail(FailureCode::Conflict, "That key is already cached."));
311 };
312 Ok(Outcome::Ok(CacheReservation { id, object, number: inserted.number as u64, upload: None, blob: None }))
313 }
314
315 /// `cache_commit`: the entry is ready; entries past the quota are
316 /// evicted, restored longest ago first.
317 pub async fn cache_commit(&self, a: CacheCommitArgs) -> Result<Outcome<CacheCommitted>> {
318 let job = check!(self.cache_job(&a.job, &a.token).await?);
319 if a.size > CACHE_MAX_ENTRY_BYTES {
320 return Ok(fail(FailureCode::Invalid, "That entry is larger than a cache entry may be."));
321 }
322 let ready = self
323 .db
324 .prepare("UPDATE cache_entries SET status = 'ready', size = ?, last_used_at = ? WHERE id = ? AND repo_id = ? AND status = 'pending' RETURNING id")
325 .bind(&[(a.size as f64).into(), rfc3339(now_ms()).into(), a.id.as_str().into(), job.repo_id.as_str().into()])?
326 .first::<Value>(None)
327 .await?;
328 if ready.is_none() {
329 return Ok(fail(FailureCode::NotFound, "No upload of that entry is in progress."));
330 }
331 #[derive(Deserialize)]
332 struct Held {
333 id: String,
334 object: String,
335 size: f64,
336 }
337 let held = self
338 .db
339 .prepare("SELECT id, object, size FROM cache_entries WHERE repo_id = ? AND status = 'ready' ORDER BY last_used_at DESC")
340 .bind(&[job.repo_id.as_str().into()])?
341 .all()
342 .await?
343 .results::<Held>()?;
344 let sizes: Vec<(String, u64)> = held.iter().map(|h| (h.id.clone(), h.size as u64)).collect();
345 let evict = to_evict(&sizes, &a.id, CACHE_REPO_QUOTA_BYTES);
346 let mut evicted = Vec::new();
347 for id in &evict {
348 if let Some(entry) = held.iter().find(|h| &h.id == id) {
349 self.db.prepare("DELETE FROM cache_entries WHERE id = ?").bind(&[id.as_str().into()])?.run().await?;
350 evicted.push(entry.object.clone());
351 }
352 }
353 Ok(Outcome::Ok(CacheCommitted { evicted }))
354 }
355
356 /// `cache_abort`.
357 pub async fn cache_abort(&self, a: CacheAbortArgs) -> Result<Outcome<bool>> {
358 let job = check!(self.cache_job(&a.job, &a.token).await?);
359 self.db
360 .prepare("DELETE FROM cache_entries WHERE id = ? AND repo_id = ? AND status = 'pending'")
361 .bind(&[a.id.as_str().into(), job.repo_id.as_str().into()])?
362 .run()
363 .await?;
364 Ok(Outcome::Ok(true))
365 }
366
367 /// Hourly: deletes entries unused for a week, saved too long ago, or
368 /// left unfinished, and their objects; then what each workspace's
369 /// cache holds today, and this month's storage, for billing.
370 pub async fn sweep_cache(&self, now: u64) -> Result<()> {
371 let at = |ms: u64| rfc3339(now.saturating_sub(ms));
372 self.db
373 .prepare(
374 "UPDATE cache_entries SET status = 'expired'
375 WHERE (status = 'ready' AND (last_used_at < ?1 OR created_at < ?2)) OR (status = 'pending' AND created_at < ?3)",
376 )
377 .bind(&[at(CACHE_UNUSED_DAYS * DAY_MS).into(), at(CACHE_MAX_AGE_DAYS * DAY_MS).into(), at(PENDING_MS).into()])?
378 .run()
379 .await?;
380 #[derive(Deserialize)]
381 struct Gone {
382 id: String,
383 object: String,
384 }
385 let gone = self
386 .db
387 .prepare("SELECT id, object FROM cache_entries WHERE status = 'expired' LIMIT 500")
388 .all()
389 .await?
390 .results::<Gone>()?;
391 if let Some(bucket) = &self.cache {
392 for chunk in gone.chunks(100) {
393 if let Err(error) = bucket.delete_multiple(chunk.iter().map(|g| g.object.as_str()).collect()).await {
394 worker::console_error!("actions: cache objects not deleted: {error}");
395 return Ok(());
396 }
397 for entry in chunk {
398 self.db.prepare("DELETE FROM cache_entries WHERE id = ?").bind(&[entry.id.as_str().into()])?.run().await?;
399 }
400 }
401 }
402 self.measure_cache(now).await
403 }
404
405 /// What each workspace's cache holds today, and this month's storage so
406 /// far reported to billing, which charges it to workspaces on the plan
407 /// once the month is over.
408 async fn measure_cache(&self, now: u64) -> Result<()> {
409 #[derive(Deserialize)]
410 struct Held {
411 namespace: String,
412 bytes: f64,
413 }
414 let today = rfc3339(now);
415 let (day, month) = (&today[..10], &today[..7]);
416 let held = self
417 .db
418 // Artifacts are kept in the same bucket, at the same price.
419 .prepare(
420 "SELECT namespace, SUM(size) AS bytes FROM (
421 SELECT namespace, size FROM cache_entries WHERE status = 'ready'
422 UNION ALL SELECT namespace, size FROM artifacts WHERE status = 'ready'
423 ) GROUP BY namespace",
424 )
425 .all()
426 .await?
427 .results::<Held>()?;
428 for workspace in &held {
429 self.db
430 .prepare(
431 "INSERT INTO cache_days (namespace, day, bytes) VALUES (?1, ?2, ?3)
432 ON CONFLICT (namespace, day) DO UPDATE SET bytes = max(bytes, ?3)",
433 )
434 .bind(&[workspace.namespace.as_str().into(), day.into(), workspace.bytes.into()])?
435 .run()
436 .await?;
437 }
438 #[derive(Deserialize)]
439 struct Day {
440 namespace: String,
441 bytes: f64,
442 }
443 let days = self
444 .db
445 .prepare("SELECT namespace, bytes FROM cache_days WHERE substr(day, 1, 7) = ?")
446 .bind(&[month.into()])?
447 .all()
448 .await?
449 .results::<Day>()?;
450 let mut by_workspace: std::collections::BTreeMap<String, Vec<u64>> = std::collections::BTreeMap::new();
451 for d in days {
452 by_workspace.entry(d.namespace).or_default().push(d.bytes as u64);
453 }
454 for (workspace, days) in by_workspace {
455 let months = gb_months(&days);
456 let cost = storage_cost(months);
457 if cost <= 0 {
458 continue;
459 }
460 let noted: Result<bool> = g1t_kit::call(
461 &self.billing,
462 "note_pending",
463 &NotePendingArgs { workspace: workspace.clone(), source: "cache".to_owned(), cost_micros: cost, detail: Some(format!("{months:.2} GB-months")) },
464 )
465 .await;
466 if let Err(error) = noted {
467 worker::console_error!("actions: cache storage not reported for {workspace}: {error}");
468 }
469 }
470 Ok(())
471 }
472}
473
474#[cfg(test)]
475mod tests {
476 use super::*;
477
478 fn entries(sizes: &[(&str, u64)]) -> Vec<(String, u64)> {
479 sizes.iter().map(|(id, size)| ((*id).to_owned(), *size)).collect()
480 }
481
482 #[test]
483 fn a_run_restores_from_its_ref_then_its_base_then_the_default_branch() {
484 let (restore, save) = scopes("refs/heads/feature", None, "main", true);
485 assert_eq!(restore, ["refs/heads/feature", "refs/heads/main"]);
486 assert_eq!(save, "refs/heads/feature");
487 // A pull request: its own merge ref, the branch it merges into, the default.
488 let (restore, save) = scopes("refs/pull/7/merge", Some("release/1.x"), "main", true);
489 assert_eq!(restore, ["refs/pull/7/merge", "refs/heads/release/1.x", "refs/heads/main"]);
490 assert_eq!(save, "refs/pull/7/merge");
491 // The default branch reads only its own.
492 let (restore, _) = scopes("refs/heads/main", Some("main"), "main", true);
493 assert_eq!(restore, ["refs/heads/main"]);
494 // From outside: it saves where nothing else reads.
495 let (restore, save) = scopes("refs/pull/9/merge", Some("main"), "main", false);
496 assert_eq!(save, "untrusted:refs/pull/9/merge");
497 assert_eq!(restore, ["untrusted:refs/pull/9/merge", "refs/heads/main"]);
498 for trusted_ref in ["refs/heads/main", "refs/pull/9/merge", "refs/heads/feature"] {
499 let (others, _) = scopes(trusted_ref, Some("main"), "main", true);
500 assert!(!others.contains(&save), "{trusted_ref} must never read an untrusted entry");
501 }
502 }
503
504 #[test]
505 fn eviction_takes_the_entries_restored_longest_ago() {
506 // Newest use first.
507 let held = entries(&[("new", 4), ("b", 3), ("c", 3), ("old", 2)]);
508 assert_eq!(to_evict(&held, "new", 12), Vec::<String>::new());
509 assert_eq!(to_evict(&held, "new", 10), ["old"]);
510 assert_eq!(to_evict(&held, "new", 7), ["old", "c"]);
511 // The entry just saved stays, even when it alone is past the quota.
512 assert_eq!(to_evict(&entries(&[("big", 20), ("a", 1)]), "big", 10), ["a"]);
513 }
514
515 #[test]
516 fn keys_are_checked() {
517 assert!(valid_key("cargo-Linux-abc123"));
518 assert!(!valid_key(""));
519 assert!(!valid_key("a,b"));
520 assert!(!valid_key(&"k".repeat(513)));
521 }
522
523 #[test]
524 fn storage_is_charged_by_the_gb_month_at_r2s_price() {
525 // 10 GB held for 30 days is 10 GB-months: $0.15.
526 let month = vec![10_000_000_000u64; 30];
527 assert!((gb_months(&month) - 10.0).abs() < 1e-9);
528 assert_eq!(storage_cost(gb_months(&month)), 150_000);
529 // A day of 1 GB: a thirtieth of a GB-month, rounded up.
530 assert_eq!(storage_cost(gb_months(&[1_000_000_000])), 500);
531 assert_eq!(storage_cost(0.0), 0);
532 }
533}