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 incidents | 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 | //! | |
| Merge main into the run-protection branch | 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 | |
| Actions: keep workflow runs safe | 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 | |
| Merge main into the run-protection branch | 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 | |
| Actions: keep workflow runs safe | 16 | //! starts with one of its restore keys. |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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 | |
| Merge main into the run-protection branch | 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. | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 25 | //! |
| Merge main into the run-protection branch | 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. | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 30 | |
| 31 | use g1t_contracts::FailureCode; | |
| 32 | use g1t_contracts::Outcome; | |
| 33 | use g1t_contracts::actions::{ | |
| 34 | CACHE_MAX_AGE_DAYS, CACHE_MAX_ENTRY_BYTES, CACHE_MICROS_PER_GB_MONTH, CACHE_REPO_QUOTA_BYTES, CACHE_UNUSED_DAYS, | |
| Merge main into the run-protection branch | 35 | CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs, CacheUploadArgs, |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 36 | }; |
| 37 | use g1t_contracts::billing::NotePendingArgs; | |
| 38 | use g1t_contracts::new_id; | |
| 39 | use g1t_contracts::time::rfc3339; | |
| 40 | use g1t_kit::now_ms; | |
| 41 | use serde::Deserialize; | |
| 42 | use serde_json::Value; | |
| 43 | use worker::Result; | |
| 44 | ||
| 45 | use crate::{Actions, check, fail}; | |
| 46 | ||
| 47 | const DAY_MS: u64 = 24 * 60 * 60 * 1000; | |
| 48 | /// An upload not finished after this long is given up. | |
| 49 | const PENDING_MS: u64 = 6 * 60 * 60 * 1000; | |
| 50 | /// A GB, as storage is billed. | |
| 51 | const GB: f64 = 1_000_000_000.0; | |
| 52 | ||
| 53 | #[derive(Debug, Deserialize)] | |
| 54 | struct EntryRow { | |
| 55 | id: String, | |
| 56 | key: String, | |
| 57 | object: String, | |
| 58 | size: f64, | |
| Merge main into the run-protection branch | 59 | created_at: String, |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 60 | } |
| 61 | ||
| 62 | /// The longest key: GitHub's limit. | |
| 63 | const MAX_KEY_CHARS: usize = 512; | |
| 64 | ||
| 65 | /// Whether a key can be kept: 1 to 512 characters, no commas (GitHub's rule). | |
| 66 | pub(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. | |
| 73 | pub(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. | |
| 91 | pub(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. | |
| 96 | pub(crate) fn storage_cost(gb_months: f64) -> i64 { | |
| 97 | (gb_months * CACHE_MICROS_PER_GB_MONTH as f64).ceil() as i64 | |
| 98 | } | |
| 99 | ||
| Merge main into the run-protection branch | 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. | |
| 102 | pub(crate) struct CacheScope { | |
| 103 | pub repo_id: String, | |
| 104 | pub restore: Vec<String>, | |
| 105 | pub save: String, | |
| 106 | } | |
| 107 | ||
| Actions: keep workflow runs safe | 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. | |
| 112 | pub(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 | ||
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 124 | impl Actions { |
| Merge main into the run-protection branch | 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 }) | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 136 | } |
| 137 | ||
| Merge main into the run-protection branch | 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 | |
| Actions: keep workflow runs safe | 141 | } |
| 142 | ||
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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?); | |
| Merge main into the run-protection branch | 146 | let scope = self.cache_scope(&job).await?; |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 147 | let now = now_ms(); |
| 148 | let fresh = rfc3339(now.saturating_sub(CACHE_MAX_AGE_DAYS * DAY_MS)); | |
| Merge main into the run-protection branch | 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(); | |
| Actions: keep workflow runs safe | 152 | let mut found = None; |
| Merge main into the run-protection branch | 153 | 'scopes: for ref_scope in &scope.restore { |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 154 | found = self |
| 155 | .db | |
| 156 | .prepare( | |
| Merge main into the run-protection branch | 157 | "SELECT id, key, object, size, created_at FROM cache_entries |
| Actions: keep workflow runs safe | 158 | 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 incidents | 159 | ) |
| Merge main into the run-protection branch | 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 | ])? | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 167 | .first::<EntryRow>(None) |
| 168 | .await?; | |
| Actions: keep workflow runs safe | 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( | |
| Merge main into the run-protection branch | 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 | |
| Actions: keep workflow runs safe | 178 | AND substr(key, 1, length(?2)) = ?2 |
| 179 | ORDER BY created_at DESC LIMIT 1", | |
| 180 | ) | |
| Merge main into the run-protection branch | 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 | ])? | |
| Actions: keep workflow runs safe | 188 | .first::<EntryRow>(None) |
| 189 | .await?; | |
| 190 | if found.is_some() { | |
| 191 | break 'scopes; | |
| 192 | } | |
| 193 | } | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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?; | |
| Merge main into the run-protection branch | 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 }))) | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 203 | } |
| 204 | ||
| Merge main into the run-protection branch | 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 | ||
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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?); | |
| Merge main into the run-protection branch | 250 | let scope = self.cache_scope(&job).await?; |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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); | |
| Merge main into the run-protection branch | 262 | let version = a.version.clone().unwrap_or_default(); |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 263 | // An upload left unfinished long ago no longer holds its key. |
| 264 | self.db | |
| Actions: keep workflow runs safe | 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(&[ | |
| Merge main into the run-protection branch | 270 | scope.repo_id.as_str().into(), |
| 271 | scope.save.as_str().into(), | |
| Actions: keep workflow runs safe | 272 | a.key.as_str().into(), |
| Merge main into the run-protection branch | 273 | version.as_str().into(), |
| Actions: keep workflow runs safe | 274 | rfc3339(now.saturating_sub(PENDING_MS)).into(), |
| 275 | ])? | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 276 | .run() |
| 277 | .await?; | |
| 278 | let id = new_id("cache", now); | |
| Merge main into the run-protection branch | 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. | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 286 | let inserted = self |
| 287 | .db | |
| 288 | .prepare( | |
| Merge main into the run-protection branch | 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) | |
| Actions: keep workflow runs safe | 291 | ON CONFLICT (repo_id, scope, key, version) DO UPDATE SET |
| Merge main into the run-protection branch | 292 | 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 incidents | 293 | WHERE cache_entries.status = 'expired' |
| Merge main into the run-protection branch | 294 | RETURNING rowid AS number", |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 295 | ) |
| 296 | .bind(&[ | |
| 297 | id.as_str().into(), | |
| Merge main into the run-protection branch | 298 | scope.repo_id.as_str().into(), |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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(), | |
| Merge main into the run-protection branch | 304 | version.as_str().into(), |
| 305 | scope.save.as_str().into(), | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 306 | ])? |
| Merge main into the run-protection branch | 307 | .first::<Inserted>(None) |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 308 | .await?; |
| Merge main into the run-protection branch | 309 | let Some(inserted) = inserted else { |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 310 | return Ok(fail(FailureCode::Conflict, "That key is already cached.")); |
| Merge main into the run-protection branch | 311 | }; |
| 312 | 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 incidents | 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 | |
| Merge main into the run-protection branch | 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 | ) | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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)] | |
| 475 | mod 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] | |
| Actions: keep workflow runs safe | 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] | |
| Merge main into the run-protection branch | 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] | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.