Skip to content
446 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Fast pages, required checks on the branch, self-hosted runners, honest incidents1//! 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//! - A key is written once. A restore finds its key exactly, else the
5//! newest entry whose key starts with one of its restore keys.
6//! - An entry is at most `CACHE_MAX_ENTRY_BYTES`. A repository's entries
7//! hold at most `CACHE_REPO_QUOTA_BYTES` together: saving past it evicts
8//! the entries restored longest ago.
9//! - An entry not restored for `CACHE_UNUSED_DAYS`, or saved more than
10//! `CACHE_MAX_AGE_DAYS` ago, is deleted by the hourly sweep, which also
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R211//! writes down what each workspace holds, its artifacts included (they
12//! are in the same bucket), and reports this month's storage to billing
13//! (`note_pending`, source `cache`) at R2's price.
Fast pages, required checks on the branch, self-hosted runners, honest incidents14//!
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R215//! The API calls these with the job's token, or its runtime token for the
16//! toolkit's protocols (runtime.rs), which is checked here. An entry the
17//! toolkit saves has the toolkit's version, and is found only by the same
18//! version; g1t's own `actions/cache` saves and finds entries without one.
Fast pages, required checks on the branch, self-hosted runners, honest incidents19
20use g1t_contracts::FailureCode;
21use g1t_contracts::Outcome;
22use g1t_contracts::actions::{
23 CACHE_MAX_AGE_DAYS, CACHE_MAX_ENTRY_BYTES, CACHE_MICROS_PER_GB_MONTH, CACHE_REPO_QUOTA_BYTES, CACHE_UNUSED_DAYS,
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R224 CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs, CacheUploadArgs,
Fast pages, required checks on the branch, self-hosted runners, honest incidents25};
26use g1t_contracts::billing::NotePendingArgs;
27use g1t_contracts::new_id;
28use g1t_contracts::time::rfc3339;
29use g1t_kit::now_ms;
30use serde::Deserialize;
31use serde_json::Value;
32use worker::Result;
33
34use crate::{Actions, check, fail};
35
36const DAY_MS: u64 = 24 * 60 * 60 * 1000;
37/// An upload not finished after this long is given up.
38const PENDING_MS: u64 = 6 * 60 * 60 * 1000;
39/// A GB, as storage is billed.
40const GB: f64 = 1_000_000_000.0;
41
42#[derive(Debug, Deserialize)]
43struct EntryRow {
44 id: String,
45 key: String,
46 object: String,
47 size: f64,
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R248 created_at: String,
Fast pages, required checks on the branch, self-hosted runners, honest incidents49}
50
51/// The longest key: GitHub's limit.
52const MAX_KEY_CHARS: usize = 512;
53
54/// Whether a key can be kept: 1 to 512 characters, no commas (GitHub's rule).
55pub(crate) fn valid_key(key: &str) -> bool {
56 !key.is_empty() && key.chars().count() <= MAX_KEY_CHARS && !key.contains(',')
57}
58
59/// Which of `entries` (id, size, newest use first) to evict so that the
60/// repository holds at most `quota` with `keep` (just saved) among them.
61/// The entry just saved is never evicted.
62pub(crate) fn to_evict(entries: &[(String, u64)], keep: &str, quota: u64) -> Vec<String> {
63 let mut total: u64 = entries.iter().map(|(_, size)| size).sum();
64 let mut out = Vec::new();
65 for (id, size) in entries.iter().rev() {
66 if total <= quota {
67 break;
68 }
69 if id == keep {
70 continue;
71 }
72 total = total.saturating_sub(*size);
73 out.push(id.clone());
74 }
75 out
76}
77
78/// A month's GB-months from the bytes held each day so far: each day's
79/// bytes over 30 days.
80pub(crate) fn gb_months(days: &[u64]) -> f64 {
81 days.iter().map(|bytes| *bytes as f64).sum::<f64>() / GB / 30.0
82}
83
84/// What `gb_months` of cache cost g1t, in millionths of a dollar.
85pub(crate) fn storage_cost(gb_months: f64) -> i64 {
86 (gb_months * CACHE_MICROS_PER_GB_MONTH as f64).ceil() as i64
87}
88
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R289/// Where a job's cache entries are found and saved. Today that is its
90/// repository: every job of a repository restores what any of its jobs
91/// saved. Narrowing that (by branch, say) belongs here, so that g1t's own
92/// `actions/cache` and the toolkit's protocols follow the same rule.
93pub(crate) struct CacheScope {
94 pub repo_id: String,
95}
96
97pub(crate) fn cache_scope(job: &crate::plan::JobRow) -> CacheScope {
98 CacheScope { repo_id: job.repo_id.clone() }
99}
100
Fast pages, required checks on the branch, self-hosted runners, honest incidents101impl Actions {
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2102 /// A job by its own token or its runtime token (runtime.rs).
Fast pages, required checks on the branch, self-hosted runners, honest incidents103 async fn cache_job(&self, job: &str, token: &str) -> Result<Outcome<crate::plan::JobRow>> {
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2104 self.job_for_credential(job, token).await
Fast pages, required checks on the branch, self-hosted runners, honest incidents105 }
106
107 /// `cache_lookup`.
108 pub async fn cache_lookup(&self, a: CacheLookupArgs) -> Result<Outcome<Option<CacheHit>>> {
109 let job = check!(self.cache_job(&a.job, &a.token).await?);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2110 let scope = cache_scope(&job);
Fast pages, required checks on the branch, self-hosted runners, honest incidents111 let now = now_ms();
112 let fresh = rfc3339(now.saturating_sub(CACHE_MAX_AGE_DAYS * DAY_MS));
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2113 // g1t's own `actions/cache` has no version; the toolkit's client
114 // finds only entries of its own version.
115 let version = a.version.clone().unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents116 let exact = self
117 .db
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2118 .prepare(
119 "SELECT id, key, object, size, created_at FROM cache_entries
120 WHERE repo_id = ? AND key = ? AND version = ? AND status = 'ready' AND created_at > ?",
121 )
122 .bind(&[scope.repo_id.as_str().into(), a.key.as_str().into(), version.as_str().into(), fresh.as_str().into()])?
Fast pages, required checks on the branch, self-hosted runners, honest incidents123 .first::<EntryRow>(None)
124 .await?;
125 let mut found = exact;
126 for prefix in a.restore.iter().map(|p| p.trim()).filter(|p| !p.is_empty()) {
127 if found.is_some() {
128 break;
129 }
130 found = self
131 .db
132 .prepare(
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2133 "SELECT id, key, object, size, created_at FROM cache_entries
134 WHERE repo_id = ?1 AND version = ?4 AND status = 'ready' AND created_at > ?3 AND substr(key, 1, length(?2)) = ?2
Fast pages, required checks on the branch, self-hosted runners, honest incidents135 ORDER BY created_at DESC LIMIT 1",
136 )
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2137 .bind(&[scope.repo_id.as_str().into(), prefix.into(), fresh.as_str().into(), version.as_str().into()])?
Fast pages, required checks on the branch, self-hosted runners, honest incidents138 .first::<EntryRow>(None)
139 .await?;
140 }
141 let Some(entry) = found else { return Ok(Outcome::Ok(None)) };
142 self.db
143 .prepare("UPDATE cache_entries SET last_used_at = ? WHERE id = ?")
144 .bind(&[rfc3339(now).into(), entry.id.as_str().into()])?
145 .run()
146 .await?;
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2147 let blob = a.version.as_ref().and_then(|_| self.download_token("cache", &entry.id, &entry.object, crate::runtime::DOWNLOAD_SECONDS));
148 Ok(Outcome::Ok(Some(CacheHit { key: entry.key, object: entry.object, size: entry.size as u64, created_at: entry.created_at, blob })))
149 }
150
151 /// `cache_upload`.
152 pub async fn cache_upload(&self, a: CacheUploadArgs) -> Result<Outcome<CacheReservation>> {
153 let job = check!(self.cache_job(&a.job, &a.token).await?);
154 let scope = cache_scope(&job);
155 #[derive(Deserialize)]
156 struct Pending {
157 id: String,
158 object: String,
159 number: f64,
160 upload: Option<String>,
161 }
162 let columns = "SELECT id, object, rowid AS number, upload FROM cache_entries";
163 let pending = match (a.number, a.key.as_deref()) {
164 (Some(number), _) => {
165 self.db
166 .prepare(format!("{columns} WHERE rowid = ? AND repo_id = ? AND status = 'pending'"))
167 .bind(&[(number as f64).into(), scope.repo_id.as_str().into()])?
168 .first::<Pending>(None)
169 .await?
170 }
171 (None, Some(key)) => {
172 self.db
173 .prepare(format!("{columns} WHERE repo_id = ? AND key = ? AND version = ? AND status = 'pending'"))
174 .bind(&[scope.repo_id.as_str().into(), key.into(), a.version.clone().unwrap_or_default().into()])?
175 .first::<Pending>(None)
176 .await?
177 }
178 (None, None) => None,
179 };
180 let Some(pending) = pending else {
181 return Ok(fail(FailureCode::NotFound, "No upload of that entry is in progress."));
182 };
183 let blob = match &pending.upload {
184 Some(upload) => match self.upload_token("cache", &pending.id, &pending.object, upload) {
185 Some(blob) => Some(blob),
186 None => return Ok(fail(FailureCode::Invalid, "The toolkit's storage is not set up here: the actions service has no ACTIONS_KEY.")),
187 },
188 None => None,
189 };
190 Ok(Outcome::Ok(CacheReservation { id: pending.id, object: pending.object, number: pending.number as u64, upload: pending.upload, blob }))
Fast pages, required checks on the branch, self-hosted runners, honest incidents191 }
192
193 /// `cache_reserve`.
194 pub async fn cache_reserve(&self, a: CacheReserveArgs) -> Result<Outcome<CacheReservation>> {
195 let job = check!(self.cache_job(&a.job, &a.token).await?);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2196 let scope = cache_scope(&job);
Fast pages, required checks on the branch, self-hosted runners, honest incidents197 if !valid_key(&a.key) {
198 return Ok(fail(FailureCode::Invalid, format!("A cache key is 1 to {MAX_KEY_CHARS} characters, without commas.")));
199 }
200 if a.size > CACHE_MAX_ENTRY_BYTES {
201 return Ok(fail(
202 FailureCode::Invalid,
203 format!("It is {} MB; a cache entry is at most {} MB.", a.size / 1_048_576, CACHE_MAX_ENTRY_BYTES / 1_048_576),
204 ));
205 }
206 let now = now_ms();
207 let at = rfc3339(now);
208 // An upload left unfinished long ago no longer holds its key.
209 self.db
210 .prepare("UPDATE cache_entries SET status = 'expired' WHERE repo_id = ? AND key = ? AND status = 'pending' AND created_at < ?")
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2211 .bind(&[scope.repo_id.as_str().into(), a.key.as_str().into(), rfc3339(now.saturating_sub(PENDING_MS)).into()])?
Fast pages, required checks on the branch, self-hosted runners, honest incidents212 .run()
213 .await?;
214 let id = new_id("cache", now);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2215 let object = format!("c/{}/{id}", scope.repo_id);
216 let version = a.version.clone().unwrap_or_default();
217 #[derive(Deserialize)]
218 struct Inserted {
219 number: f64,
220 }
221 // A key is written once, whatever its version.
Fast pages, required checks on the branch, self-hosted runners, honest incidents222 let inserted = self
223 .db
224 .prepare(
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2225 "INSERT INTO cache_entries (id, repo_id, namespace, key, object, size, status, created_at, last_used_at, version)
226 VALUES (?1, ?2, ?3, ?4, ?5, ?6, 'pending', ?7, ?7, ?8)
Fast pages, required checks on the branch, self-hosted runners, honest incidents227 ON CONFLICT (repo_id, key) DO UPDATE SET
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2228 id = ?1, object = ?5, size = ?6, status = 'pending', created_at = ?7, last_used_at = ?7, version = ?8, upload = NULL
Fast pages, required checks on the branch, self-hosted runners, honest incidents229 WHERE cache_entries.status = 'expired'
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2230 RETURNING rowid AS number",
Fast pages, required checks on the branch, self-hosted runners, honest incidents231 )
232 .bind(&[
233 id.as_str().into(),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2234 scope.repo_id.as_str().into(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents235 job.namespace.as_str().into(),
236 a.key.as_str().into(),
237 object.as_str().into(),
238 (a.size as f64).into(),
239 at.as_str().into(),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2240 version.as_str().into(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents241 ])?
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2242 .first::<Inserted>(None)
Fast pages, required checks on the branch, self-hosted runners, honest incidents243 .await?;
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2244 let Some(inserted) = inserted else {
Fast pages, required checks on the branch, self-hosted runners, honest incidents245 return Ok(fail(FailureCode::Conflict, "That key is already cached."));
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2246 };
247 Ok(Outcome::Ok(CacheReservation { id, object, number: inserted.number as u64, upload: None, blob: None }))
Fast pages, required checks on the branch, self-hosted runners, honest incidents248 }
249
250 /// `cache_commit`: the entry is ready; entries past the quota are
251 /// evicted, restored longest ago first.
252 pub async fn cache_commit(&self, a: CacheCommitArgs) -> Result<Outcome<CacheCommitted>> {
253 let job = check!(self.cache_job(&a.job, &a.token).await?);
254 if a.size > CACHE_MAX_ENTRY_BYTES {
255 return Ok(fail(FailureCode::Invalid, "That entry is larger than a cache entry may be."));
256 }
257 let ready = self
258 .db
259 .prepare("UPDATE cache_entries SET status = 'ready', size = ?, last_used_at = ? WHERE id = ? AND repo_id = ? AND status = 'pending' RETURNING id")
260 .bind(&[(a.size as f64).into(), rfc3339(now_ms()).into(), a.id.as_str().into(), job.repo_id.as_str().into()])?
261 .first::<Value>(None)
262 .await?;
263 if ready.is_none() {
264 return Ok(fail(FailureCode::NotFound, "No upload of that entry is in progress."));
265 }
266 #[derive(Deserialize)]
267 struct Held {
268 id: String,
269 object: String,
270 size: f64,
271 }
272 let held = self
273 .db
274 .prepare("SELECT id, object, size FROM cache_entries WHERE repo_id = ? AND status = 'ready' ORDER BY last_used_at DESC")
275 .bind(&[job.repo_id.as_str().into()])?
276 .all()
277 .await?
278 .results::<Held>()?;
279 let sizes: Vec<(String, u64)> = held.iter().map(|h| (h.id.clone(), h.size as u64)).collect();
280 let evict = to_evict(&sizes, &a.id, CACHE_REPO_QUOTA_BYTES);
281 let mut evicted = Vec::new();
282 for id in &evict {
283 if let Some(entry) = held.iter().find(|h| &h.id == id) {
284 self.db.prepare("DELETE FROM cache_entries WHERE id = ?").bind(&[id.as_str().into()])?.run().await?;
285 evicted.push(entry.object.clone());
286 }
287 }
288 Ok(Outcome::Ok(CacheCommitted { evicted }))
289 }
290
291 /// `cache_abort`.
292 pub async fn cache_abort(&self, a: CacheAbortArgs) -> Result<Outcome<bool>> {
293 let job = check!(self.cache_job(&a.job, &a.token).await?);
294 self.db
295 .prepare("DELETE FROM cache_entries WHERE id = ? AND repo_id = ? AND status = 'pending'")
296 .bind(&[a.id.as_str().into(), job.repo_id.as_str().into()])?
297 .run()
298 .await?;
299 Ok(Outcome::Ok(true))
300 }
301
302 /// Hourly: deletes entries unused for a week, saved too long ago, or
303 /// left unfinished, and their objects; then what each workspace's
304 /// cache holds today, and this month's storage, for billing.
305 pub async fn sweep_cache(&self, now: u64) -> Result<()> {
306 let at = |ms: u64| rfc3339(now.saturating_sub(ms));
307 self.db
308 .prepare(
309 "UPDATE cache_entries SET status = 'expired'
310 WHERE (status = 'ready' AND (last_used_at < ?1 OR created_at < ?2)) OR (status = 'pending' AND created_at < ?3)",
311 )
312 .bind(&[at(CACHE_UNUSED_DAYS * DAY_MS).into(), at(CACHE_MAX_AGE_DAYS * DAY_MS).into(), at(PENDING_MS).into()])?
313 .run()
314 .await?;
315 #[derive(Deserialize)]
316 struct Gone {
317 id: String,
318 object: String,
319 }
320 let gone = self
321 .db
322 .prepare("SELECT id, object FROM cache_entries WHERE status = 'expired' LIMIT 500")
323 .all()
324 .await?
325 .results::<Gone>()?;
326 if let Some(bucket) = &self.cache {
327 for chunk in gone.chunks(100) {
328 if let Err(error) = bucket.delete_multiple(chunk.iter().map(|g| g.object.as_str()).collect()).await {
329 worker::console_error!("actions: cache objects not deleted: {error}");
330 return Ok(());
331 }
332 for entry in chunk {
333 self.db.prepare("DELETE FROM cache_entries WHERE id = ?").bind(&[entry.id.as_str().into()])?.run().await?;
334 }
335 }
336 }
337 self.measure_cache(now).await
338 }
339
340 /// What each workspace's cache holds today, and this month's storage so
341 /// far reported to billing, which charges it to workspaces on the plan
342 /// once the month is over.
343 async fn measure_cache(&self, now: u64) -> Result<()> {
344 #[derive(Deserialize)]
345 struct Held {
346 namespace: String,
347 bytes: f64,
348 }
349 let today = rfc3339(now);
350 let (day, month) = (&today[..10], &today[..7]);
351 let held = self
352 .db
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2353 // Artifacts are kept in the same bucket, at the same price.
354 .prepare(
355 "SELECT namespace, SUM(size) AS bytes FROM (
356 SELECT namespace, size FROM cache_entries WHERE status = 'ready'
357 UNION ALL SELECT namespace, size FROM artifacts WHERE status = 'ready'
358 ) GROUP BY namespace",
359 )
Fast pages, required checks on the branch, self-hosted runners, honest incidents360 .all()
361 .await?
362 .results::<Held>()?;
363 for workspace in &held {
364 self.db
365 .prepare(
366 "INSERT INTO cache_days (namespace, day, bytes) VALUES (?1, ?2, ?3)
367 ON CONFLICT (namespace, day) DO UPDATE SET bytes = max(bytes, ?3)",
368 )
369 .bind(&[workspace.namespace.as_str().into(), day.into(), workspace.bytes.into()])?
370 .run()
371 .await?;
372 }
373 #[derive(Deserialize)]
374 struct Day {
375 namespace: String,
376 bytes: f64,
377 }
378 let days = self
379 .db
380 .prepare("SELECT namespace, bytes FROM cache_days WHERE substr(day, 1, 7) = ?")
381 .bind(&[month.into()])?
382 .all()
383 .await?
384 .results::<Day>()?;
385 let mut by_workspace: std::collections::BTreeMap<String, Vec<u64>> = std::collections::BTreeMap::new();
386 for d in days {
387 by_workspace.entry(d.namespace).or_default().push(d.bytes as u64);
388 }
389 for (workspace, days) in by_workspace {
390 let months = gb_months(&days);
391 let cost = storage_cost(months);
392 if cost <= 0 {
393 continue;
394 }
395 let noted: Result<bool> = g1t_kit::call(
396 &self.billing,
397 "note_pending",
398 &NotePendingArgs { workspace: workspace.clone(), source: "cache".to_owned(), cost_micros: cost, detail: Some(format!("{months:.2} GB-months")) },
399 )
400 .await;
401 if let Err(error) = noted {
402 worker::console_error!("actions: cache storage not reported for {workspace}: {error}");
403 }
404 }
405 Ok(())
406 }
407}
408
409#[cfg(test)]
410mod tests {
411 use super::*;
412
413 fn entries(sizes: &[(&str, u64)]) -> Vec<(String, u64)> {
414 sizes.iter().map(|(id, size)| ((*id).to_owned(), *size)).collect()
415 }
416
417 #[test]
418 fn eviction_takes_the_entries_restored_longest_ago() {
419 // Newest use first.
420 let held = entries(&[("new", 4), ("b", 3), ("c", 3), ("old", 2)]);
421 assert_eq!(to_evict(&held, "new", 12), Vec::<String>::new());
422 assert_eq!(to_evict(&held, "new", 10), ["old"]);
423 assert_eq!(to_evict(&held, "new", 7), ["old", "c"]);
424 // The entry just saved stays, even when it alone is past the quota.
425 assert_eq!(to_evict(&entries(&[("big", 20), ("a", 1)]), "big", 10), ["a"]);
426 }
427
428 #[test]
429 fn keys_are_checked() {
430 assert!(valid_key("cargo-Linux-abc123"));
431 assert!(!valid_key(""));
432 assert!(!valid_key("a,b"));
433 assert!(!valid_key(&"k".repeat(513)));
434 }
435
436 #[test]
437 fn storage_is_charged_by_the_gb_month_at_r2s_price() {
438 // 10 GB held for 30 days is 10 GB-months: $0.15.
439 let month = vec![10_000_000_000u64; 30];
440 assert!((gb_months(&month) - 10.0).abs() < 1e-9);
441 assert_eq!(storage_cost(gb_months(&month)), 150_000);
442 // A day of 1 GB: a thirtieth of a GB-month, rounded up.
443 assert_eq!(storage_cost(gb_months(&[1_000_000_000])), 500);
444 assert_eq!(storage_cost(0.0), 0);
445 }
446}