flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/actions/src/cache.rs

364 lines15,104 bytesCodeBlame
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//! - 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
11//! writes down what each workspace holds and reports this month's storage
12//! to billing (`note_pending`, source `cache`) at R2's price.
13//!
14//! The API calls these with the job's token, which is checked here.
15
16use g1t_contracts::FailureCode;
17use g1t_contracts::Outcome;
18use g1t_contracts::actions::{
19 CACHE_MAX_AGE_DAYS, CACHE_MAX_ENTRY_BYTES, CACHE_MICROS_PER_GB_MONTH, CACHE_REPO_QUOTA_BYTES, CACHE_UNUSED_DAYS,
20 CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs, JobCallArgs,
21};
22use g1t_contracts::billing::NotePendingArgs;
23use g1t_contracts::new_id;
24use g1t_contracts::time::rfc3339;
25use g1t_kit::now_ms;
26use serde::Deserialize;
27use serde_json::Value;
28use worker::Result;
29
30use crate::{Actions, check, fail};
31
32const DAY_MS: u64 = 24 * 60 * 60 * 1000;
33/// An upload not finished after this long is given up.
34const PENDING_MS: u64 = 6 * 60 * 60 * 1000;
35/// A GB, as storage is billed.
36const GB: f64 = 1_000_000_000.0;
37
38#[derive(Debug, Deserialize)]
39struct EntryRow {
40 id: String,
41 key: String,
42 object: String,
43 size: f64,
44}
45
46/// The longest key: GitHub's limit.
47const MAX_KEY_CHARS: usize = 512;
48
49/// Whether a key can be kept: 1 to 512 characters, no commas (GitHub's rule).
50pub(crate) fn valid_key(key: &str) -> bool {
51 !key.is_empty() && key.chars().count() <= MAX_KEY_CHARS && !key.contains(',')
52}
53
54/// Which of `entries` (id, size, newest use first) to evict so that the
55/// repository holds at most `quota` with `keep` (just saved) among them.
56/// The entry just saved is never evicted.
57pub(crate) fn to_evict(entries: &[(String, u64)], keep: &str, quota: u64) -> Vec<String> {
58 let mut total: u64 = entries.iter().map(|(_, size)| size).sum();
59 let mut out = Vec::new();
60 for (id, size) in entries.iter().rev() {
61 if total <= quota {
62 break;
63 }
64 if id == keep {
65 continue;
66 }
67 total = total.saturating_sub(*size);
68 out.push(id.clone());
69 }
70 out
71}
72
73/// A month's GB-months from the bytes held each day so far: each day's
74/// bytes over 30 days.
75pub(crate) fn gb_months(days: &[u64]) -> f64 {
76 days.iter().map(|bytes| *bytes as f64).sum::<f64>() / GB / 30.0
77}
78
79/// What `gb_months` of cache cost g1t, in millionths of a dollar.
80pub(crate) fn storage_cost(gb_months: f64) -> i64 {
81 (gb_months * CACHE_MICROS_PER_GB_MONTH as f64).ceil() as i64
82}
83
84impl Actions {
85 async fn cache_job(&self, job: &str, token: &str) -> Result<Outcome<crate::plan::JobRow>> {
86 self.job_for_token(&JobCallArgs { job: job.to_owned(), token: token.to_owned(), report: Value::Null }).await
87 }
88
89 /// `cache_lookup`.
90 pub async fn cache_lookup(&self, a: CacheLookupArgs) -> Result<Outcome<Option<CacheHit>>> {
91 let job = check!(self.cache_job(&a.job, &a.token).await?);
92 let now = now_ms();
93 let fresh = rfc3339(now.saturating_sub(CACHE_MAX_AGE_DAYS * DAY_MS));
94 let exact = self
95 .db
96 .prepare("SELECT id, key, object, size FROM cache_entries WHERE repo_id = ? AND key = ? AND status = 'ready' AND created_at > ?")
97 .bind(&[job.repo_id.as_str().into(), a.key.as_str().into(), fresh.as_str().into()])?
98 .first::<EntryRow>(None)
99 .await?;
100 let mut found = exact;
101 for prefix in a.restore.iter().map(|p| p.trim()).filter(|p| !p.is_empty()) {
102 if found.is_some() {
103 break;
104 }
105 found = self
106 .db
107 .prepare(
108 "SELECT id, key, object, size FROM cache_entries
109 WHERE repo_id = ?1 AND status = 'ready' AND created_at > ?3 AND substr(key, 1, length(?2)) = ?2
110 ORDER BY created_at DESC LIMIT 1",
111 )
112 .bind(&[job.repo_id.as_str().into(), prefix.into(), fresh.as_str().into()])?
113 .first::<EntryRow>(None)
114 .await?;
115 }
116 let Some(entry) = found else { return Ok(Outcome::Ok(None)) };
117 self.db
118 .prepare("UPDATE cache_entries SET last_used_at = ? WHERE id = ?")
119 .bind(&[rfc3339(now).into(), entry.id.as_str().into()])?
120 .run()
121 .await?;
122 Ok(Outcome::Ok(Some(CacheHit { key: entry.key, object: entry.object, size: entry.size as u64 })))
123 }
124
125 /// `cache_reserve`.
126 pub async fn cache_reserve(&self, a: CacheReserveArgs) -> Result<Outcome<CacheReservation>> {
127 let job = check!(self.cache_job(&a.job, &a.token).await?);
128 if !valid_key(&a.key) {
129 return Ok(fail(FailureCode::Invalid, format!("A cache key is 1 to {MAX_KEY_CHARS} characters, without commas.")));
130 }
131 if a.size > CACHE_MAX_ENTRY_BYTES {
132 return Ok(fail(
133 FailureCode::Invalid,
134 format!("It is {} MB; a cache entry is at most {} MB.", a.size / 1_048_576, CACHE_MAX_ENTRY_BYTES / 1_048_576),
135 ));
136 }
137 let now = now_ms();
138 let at = rfc3339(now);
139 // An upload left unfinished long ago no longer holds its key.
140 self.db
141 .prepare("UPDATE cache_entries SET status = 'expired' WHERE repo_id = ? AND key = ? AND status = 'pending' AND created_at < ?")
142 .bind(&[job.repo_id.as_str().into(), a.key.as_str().into(), rfc3339(now.saturating_sub(PENDING_MS)).into()])?
143 .run()
144 .await?;
145 let id = new_id("cache", now);
146 let object = format!("c/{}/{id}", job.repo_id);
147 let inserted = self
148 .db
149 .prepare(
150 "INSERT INTO cache_entries (id, repo_id, namespace, key, object, size, status, created_at, last_used_at)
151 VALUES (?1, ?2, ?3, ?4, ?5, ?6, 'pending', ?7, ?7)
152 ON CONFLICT (repo_id, key) DO UPDATE SET
153 id = ?1, object = ?5, size = ?6, status = 'pending', created_at = ?7, last_used_at = ?7
154 WHERE cache_entries.status = 'expired'
155 RETURNING id",
156 )
157 .bind(&[
158 id.as_str().into(),
159 job.repo_id.as_str().into(),
160 job.namespace.as_str().into(),
161 a.key.as_str().into(),
162 object.as_str().into(),
163 (a.size as f64).into(),
164 at.as_str().into(),
165 ])?
166 .first::<Value>(None)
167 .await?;
168 if inserted.is_none() {
169 return Ok(fail(FailureCode::Conflict, "That key is already cached."));
170 }
171 Ok(Outcome::Ok(CacheReservation { id, object }))
172 }
173
174 /// `cache_commit`: the entry is ready; entries past the quota are
175 /// evicted, restored longest ago first.
176 pub async fn cache_commit(&self, a: CacheCommitArgs) -> Result<Outcome<CacheCommitted>> {
177 let job = check!(self.cache_job(&a.job, &a.token).await?);
178 if a.size > CACHE_MAX_ENTRY_BYTES {
179 return Ok(fail(FailureCode::Invalid, "That entry is larger than a cache entry may be."));
180 }
181 let ready = self
182 .db
183 .prepare("UPDATE cache_entries SET status = 'ready', size = ?, last_used_at = ? WHERE id = ? AND repo_id = ? AND status = 'pending' RETURNING id")
184 .bind(&[(a.size as f64).into(), rfc3339(now_ms()).into(), a.id.as_str().into(), job.repo_id.as_str().into()])?
185 .first::<Value>(None)
186 .await?;
187 if ready.is_none() {
188 return Ok(fail(FailureCode::NotFound, "No upload of that entry is in progress."));
189 }
190 #[derive(Deserialize)]
191 struct Held {
192 id: String,
193 object: String,
194 size: f64,
195 }
196 let held = self
197 .db
198 .prepare("SELECT id, object, size FROM cache_entries WHERE repo_id = ? AND status = 'ready' ORDER BY last_used_at DESC")
199 .bind(&[job.repo_id.as_str().into()])?
200 .all()
201 .await?
202 .results::<Held>()?;
203 let sizes: Vec<(String, u64)> = held.iter().map(|h| (h.id.clone(), h.size as u64)).collect();
204 let evict = to_evict(&sizes, &a.id, CACHE_REPO_QUOTA_BYTES);
205 let mut evicted = Vec::new();
206 for id in &evict {
207 if let Some(entry) = held.iter().find(|h| &h.id == id) {
208 self.db.prepare("DELETE FROM cache_entries WHERE id = ?").bind(&[id.as_str().into()])?.run().await?;
209 evicted.push(entry.object.clone());
210 }
211 }
212 Ok(Outcome::Ok(CacheCommitted { evicted }))
213 }
214
215 /// `cache_abort`.
216 pub async fn cache_abort(&self, a: CacheAbortArgs) -> Result<Outcome<bool>> {
217 let job = check!(self.cache_job(&a.job, &a.token).await?);
218 self.db
219 .prepare("DELETE FROM cache_entries WHERE id = ? AND repo_id = ? AND status = 'pending'")
220 .bind(&[a.id.as_str().into(), job.repo_id.as_str().into()])?
221 .run()
222 .await?;
223 Ok(Outcome::Ok(true))
224 }
225
226 /// Hourly: deletes entries unused for a week, saved too long ago, or
227 /// left unfinished, and their objects; then what each workspace's
228 /// cache holds today, and this month's storage, for billing.
229 pub async fn sweep_cache(&self, now: u64) -> Result<()> {
230 let at = |ms: u64| rfc3339(now.saturating_sub(ms));
231 self.db
232 .prepare(
233 "UPDATE cache_entries SET status = 'expired'
234 WHERE (status = 'ready' AND (last_used_at < ?1 OR created_at < ?2)) OR (status = 'pending' AND created_at < ?3)",
235 )
236 .bind(&[at(CACHE_UNUSED_DAYS * DAY_MS).into(), at(CACHE_MAX_AGE_DAYS * DAY_MS).into(), at(PENDING_MS).into()])?
237 .run()
238 .await?;
239 #[derive(Deserialize)]
240 struct Gone {
241 id: String,
242 object: String,
243 }
244 let gone = self
245 .db
246 .prepare("SELECT id, object FROM cache_entries WHERE status = 'expired' LIMIT 500")
247 .all()
248 .await?
249 .results::<Gone>()?;
250 if let Some(bucket) = &self.cache {
251 for chunk in gone.chunks(100) {
252 if let Err(error) = bucket.delete_multiple(chunk.iter().map(|g| g.object.as_str()).collect()).await {
253 worker::console_error!("actions: cache objects not deleted: {error}");
254 return Ok(());
255 }
256 for entry in chunk {
257 self.db.prepare("DELETE FROM cache_entries WHERE id = ?").bind(&[entry.id.as_str().into()])?.run().await?;
258 }
259 }
260 }
261 self.measure_cache(now).await
262 }
263
264 /// What each workspace's cache holds today, and this month's storage so
265 /// far reported to billing, which charges it to workspaces on the plan
266 /// once the month is over.
267 async fn measure_cache(&self, now: u64) -> Result<()> {
268 #[derive(Deserialize)]
269 struct Held {
270 namespace: String,
271 bytes: f64,
272 }
273 let today = rfc3339(now);
274 let (day, month) = (&today[..10], &today[..7]);
275 let held = self
276 .db
277 .prepare("SELECT namespace, SUM(size) AS bytes FROM cache_entries WHERE status = 'ready' GROUP BY namespace")
278 .all()
279 .await?
280 .results::<Held>()?;
281 for workspace in &held {
282 self.db
283 .prepare(
284 "INSERT INTO cache_days (namespace, day, bytes) VALUES (?1, ?2, ?3)
285 ON CONFLICT (namespace, day) DO UPDATE SET bytes = max(bytes, ?3)",
286 )
287 .bind(&[workspace.namespace.as_str().into(), day.into(), workspace.bytes.into()])?
288 .run()
289 .await?;
290 }
291 #[derive(Deserialize)]
292 struct Day {
293 namespace: String,
294 bytes: f64,
295 }
296 let days = self
297 .db
298 .prepare("SELECT namespace, bytes FROM cache_days WHERE substr(day, 1, 7) = ?")
299 .bind(&[month.into()])?
300 .all()
301 .await?
302 .results::<Day>()?;
303 let mut by_workspace: std::collections::BTreeMap<String, Vec<u64>> = std::collections::BTreeMap::new();
304 for d in days {
305 by_workspace.entry(d.namespace).or_default().push(d.bytes as u64);
306 }
307 for (workspace, days) in by_workspace {
308 let months = gb_months(&days);
309 let cost = storage_cost(months);
310 if cost <= 0 {
311 continue;
312 }
313 let noted: Result<bool> = g1t_kit::call(
314 &self.billing,
315 "note_pending",
316 &NotePendingArgs { workspace: workspace.clone(), source: "cache".to_owned(), cost_micros: cost, detail: Some(format!("{months:.2} GB-months")) },
317 )
318 .await;
319 if let Err(error) = noted {
320 worker::console_error!("actions: cache storage not reported for {workspace}: {error}");
321 }
322 }
323 Ok(())
324 }
325}
326
327#[cfg(test)]
328mod tests {
329 use super::*;
330
331 fn entries(sizes: &[(&str, u64)]) -> Vec<(String, u64)> {
332 sizes.iter().map(|(id, size)| ((*id).to_owned(), *size)).collect()
333 }
334
335 #[test]
336 fn eviction_takes_the_entries_restored_longest_ago() {
337 // Newest use first.
338 let held = entries(&[("new", 4), ("b", 3), ("c", 3), ("old", 2)]);
339 assert_eq!(to_evict(&held, "new", 12), Vec::<String>::new());
340 assert_eq!(to_evict(&held, "new", 10), ["old"]);
341 assert_eq!(to_evict(&held, "new", 7), ["old", "c"]);
342 // The entry just saved stays, even when it alone is past the quota.
343 assert_eq!(to_evict(&entries(&[("big", 20), ("a", 1)]), "big", 10), ["a"]);
344 }
345
346 #[test]
347 fn keys_are_checked() {
348 assert!(valid_key("cargo-Linux-abc123"));
349 assert!(!valid_key(""));
350 assert!(!valid_key("a,b"));
351 assert!(!valid_key(&"k".repeat(513)));
352 }
353
354 #[test]
355 fn storage_is_charged_by_the_gb_month_at_r2s_price() {
356 // 10 GB held for 30 days is 10 GB-months: $0.15.
357 let month = vec![10_000_000_000u64; 30];
358 assert!((gb_months(&month) - 10.0).abs() < 1e-9);
359 assert_eq!(storage_cost(gb_months(&month)), 150_000);
360 // A day of 1 GB: a thirtieth of a GB-month, rounded up.
361 assert_eq!(storage_cost(gb_months(&[1_000_000_000])), 500);
362 assert_eq!(storage_cost(0.0), 0);
363 }
364}