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