| 1 | //! Artifacts of workflow runs: which there are, kept here, while the API |
| 2 | //! keeps their bytes in R2 (the ACTIONS_CACHE bucket, under `a/`). |
| 3 | //! |
| 4 | //! - A name is one artifact in a run. Uploading a name the run has is |
| 5 | //! refused unless the upload says `overwrite`, which replaces it. |
| 6 | //! - One artifact is at most `ARTIFACT_MAX_BYTES`, and a run's artifacts at |
| 7 | //! most `RUN_ARTIFACTS_MAX_BYTES` together. |
| 8 | //! - Each is kept for its `retention-days`, or the repository's setting |
| 9 | //! when it gives none, and never longer than that setting (1 to 90 days, |
| 10 | //! 14 unless changed). The hourly sweep deletes what has expired, and |
| 11 | //! uploads left unfinished. |
| 12 | //! - Ids are numbers, as GitHub's are. |
| 13 | //! |
| 14 | //! A sandbox reaches these with its job's token or its runtime token |
| 15 | //! (runtime.rs); people through the API and the run's page. |
| 16 | |
| 17 | use g1t_contracts::access::Capability; |
| 18 | use g1t_contracts::actions::{ |
| 19 | ARTIFACT_MAX_BYTES, ARTIFACT_RETENTION_DEFAULT_DAYS, ARTIFACT_RETENTION_MAX_DAYS, Artifact, ArtifactArgs, ArtifactBlob, ArtifactCommitArgs, |
| 20 | ArtifactList, ArtifactReservation, ArtifactReserveArgs, ArtifactRetention, ArtifactRetentionArgs, ArtifactsArgs, DeleteArtifactArgs, |
| 21 | JobArtifactsArgs, RUN_ARTIFACTS_MAX_BYTES, |
| 22 | }; |
| 23 | use g1t_contracts::time::{parse_rfc3339, rfc3339}; |
| 24 | use g1t_contracts::{FailureCode, Outcome, new_id}; |
| 25 | use g1t_kit::now_ms; |
| 26 | use serde::Deserialize; |
| 27 | use worker::Result; |
| 28 | use worker::wasm_bindgen::JsValue; |
| 29 | |
| 30 | use crate::runtime::{DOWNLOAD_SECONDS, LINK_SECONDS}; |
| 31 | use crate::{Actions, check, fail}; |
| 32 | |
| 33 | const DAY_MS: u64 = 24 * 60 * 60 * 1000; |
| 34 | /// An upload not finished after this long is given up. |
| 35 | const PENDING_MS: u64 = 6 * 60 * 60 * 1000; |
| 36 | const COLUMNS: &str = "artifacts.*, runs.git_ref AS git_ref, runs.sha AS sha"; |
| 37 | |
| 38 | #[derive(Debug, Deserialize)] |
| 39 | struct Row { |
| 40 | id: f64, |
| 41 | repo_id: String, |
| 42 | run_id: String, |
| 43 | job_id: String, |
| 44 | name: String, |
| 45 | object: String, |
| 46 | format: String, |
| 47 | size: f64, |
| 48 | digest: Option<String>, |
| 49 | status: String, |
| 50 | created_at: String, |
| 51 | updated_at: String, |
| 52 | expires_at: String, |
| 53 | #[serde(default)] |
| 54 | git_ref: Option<String>, |
| 55 | #[serde(default)] |
| 56 | sha: Option<String>, |
| 57 | } |
| 58 | |
| 59 | impl Row { |
| 60 | fn view(&self) -> Artifact { |
| 61 | Artifact { |
| 62 | id: self.id as u64, |
| 63 | name: self.name.clone(), |
| 64 | size: self.size as u64, |
| 65 | digest: self.digest.clone(), |
| 66 | format: self.format.clone(), |
| 67 | run_id: self.run_id.clone(), |
| 68 | job_id: self.job_id.clone(), |
| 69 | repo_id: self.repo_id.clone(), |
| 70 | expired: self.status == "expired", |
| 71 | created_at: self.created_at.clone(), |
| 72 | updated_at: self.updated_at.clone(), |
| 73 | expires_at: self.expires_at.clone(), |
| 74 | head_branch: self.git_ref.as_deref().map(|r| r.strip_prefix("refs/heads/").unwrap_or(r).to_owned()), |
| 75 | head_sha: self.sha.clone(), |
| 76 | } |
| 77 | } |
| 78 | } |
| 79 | |
| 80 | /// GitHub's rule for an artifact's name: 1 to 256 characters, none of |
| 81 | /// `" : < > | * ? \ /` or line breaks. |
| 82 | pub(crate) fn valid_name(name: &str) -> bool { |
| 83 | !name.trim().is_empty() && name.chars().count() <= 256 && !name.chars().any(|c| matches!(c, '"' | ':' | '<' | '>' | '|' | '*' | '?' | '\\' | '/' | '\r' | '\n')) |
| 84 | } |
| 85 | |
| 86 | /// The days an artifact is kept: what it asked for (as days, or as a time |
| 87 | /// to expire), else the repository's `setting`, and never more than the |
| 88 | /// setting. |
| 89 | pub(crate) fn retention(asked_days: u32, asked_until_ms: Option<u64>, setting: u32, now: u64) -> u32 { |
| 90 | let setting = setting.clamp(1, ARTIFACT_RETENTION_MAX_DAYS); |
| 91 | let days = match asked_until_ms { |
| 92 | _ if asked_days > 0 => asked_days, |
| 93 | Some(until) if until > now => (until - now).div_ceil(DAY_MS) as u32, |
| 94 | _ => setting, |
| 95 | }; |
| 96 | days.clamp(1, setting) |
| 97 | } |
| 98 | |
| 99 | /// Where an artifact is in R2: under its repository, by its id. The API's |
| 100 | /// blobs.rs names it the same way. |
| 101 | pub(crate) fn object_of(repo_id: &str, id: u64) -> String { |
| 102 | format!("a/{repo_id}/{id}") |
| 103 | } |
| 104 | |
| 105 | fn number(id: u64) -> JsValue { |
| 106 | JsValue::from_f64(id as f64) |
| 107 | } |
| 108 | |
| 109 | impl Actions { |
| 110 | /// How long the repository keeps artifacts. |
| 111 | pub(crate) async fn retention_setting(&self, repo_id: &str) -> Result<u32> { |
| 112 | #[derive(Deserialize)] |
| 113 | struct Setting { |
| 114 | artifact_retention_days: f64, |
| 115 | } |
| 116 | let setting = self |
| 117 | .db |
| 118 | .prepare("SELECT artifact_retention_days FROM repo_settings WHERE repo_id = ?") |
| 119 | .bind(&[repo_id.into()])? |
| 120 | .first::<Setting>(None) |
| 121 | .await?; |
| 122 | Ok(setting.map_or(ARTIFACT_RETENTION_DEFAULT_DAYS, |s| s.artifact_retention_days as u32)) |
| 123 | } |
| 124 | |
| 125 | async fn artifact_rows(&self, filter: &str, values: &[JsValue]) -> Result<Vec<Row>> { |
| 126 | self.db |
| 127 | .prepare(format!("SELECT {COLUMNS} FROM artifacts JOIN runs ON runs.id = artifacts.run_id WHERE {filter}")) |
| 128 | .bind(values)? |
| 129 | .all() |
| 130 | .await? |
| 131 | .results::<Row>() |
| 132 | } |
| 133 | |
| 134 | /// The bytes a run's artifacts hold, `except` one. |
| 135 | async fn run_bytes(&self, run: &str, except: u64) -> Result<u64> { |
| 136 | #[derive(Deserialize)] |
| 137 | struct Sum { |
| 138 | bytes: Option<f64>, |
| 139 | } |
| 140 | let sum = self |
| 141 | .db |
| 142 | .prepare("SELECT SUM(size) AS bytes FROM artifacts WHERE run_id = ? AND status <> 'expired' AND id <> ?") |
| 143 | .bind(&[run.into(), number(except)])? |
| 144 | .first::<Sum>(None) |
| 145 | .await?; |
| 146 | Ok(sum.and_then(|s| s.bytes).unwrap_or(0.0) as u64) |
| 147 | } |
| 148 | |
| 149 | /// Deletes artifacts' objects, then their rows. Objects that cannot be |
| 150 | /// deleted now are left to the sweep, their rows marked expired. |
| 151 | async fn forget_artifacts(&self, rows: &[(u64, String)]) -> Result<()> { |
| 152 | for (id, _) in rows { |
| 153 | self.db |
| 154 | .prepare("UPDATE artifacts SET status = 'expired', updated_at = ? WHERE id = ?") |
| 155 | .bind(&[rfc3339(now_ms()).into(), number(*id)])? |
| 156 | .run() |
| 157 | .await?; |
| 158 | } |
| 159 | let Some(bucket) = &self.cache else { return Ok(()) }; |
| 160 | for chunk in rows.chunks(100) { |
| 161 | if let Err(error) = bucket.delete_multiple(chunk.iter().map(|(_, object)| object.as_str()).collect()).await { |
| 162 | worker::console_error!("actions: artifact objects not deleted: {error}"); |
| 163 | return Ok(()); |
| 164 | } |
| 165 | for (id, _) in chunk { |
| 166 | self.db.prepare("DELETE FROM artifacts WHERE id = ?").bind(&[number(*id)])?.run().await?; |
| 167 | } |
| 168 | } |
| 169 | Ok(()) |
| 170 | } |
| 171 | |
| 172 | /// `artifact_reserve`. |
| 173 | pub async fn artifact_reserve(&self, a: ArtifactReserveArgs) -> Result<Outcome<ArtifactReservation>> { |
| 174 | let job = check!(self.job_for_credential(&a.job, &a.token).await?); |
| 175 | if !valid_name(&a.name) { |
| 176 | return Ok(fail( |
| 177 | FailureCode::Invalid, |
| 178 | format!("`{}` is not an artifact name: 1 to 256 characters, none of \" : < > | * ? \\ / or line breaks.", a.name), |
| 179 | )); |
| 180 | } |
| 181 | if a.size > ARTIFACT_MAX_BYTES { |
| 182 | return Ok(fail( |
| 183 | FailureCode::Invalid, |
| 184 | format!("It is {} MB; an artifact is at most {} MB.", a.size / 1_048_576, ARTIFACT_MAX_BYTES / 1_048_576), |
| 185 | )); |
| 186 | } |
| 187 | let held = self.run_bytes(&job.run_id, 0).await?; |
| 188 | if held + a.size > RUN_ARTIFACTS_MAX_BYTES { |
| 189 | return Ok(fail( |
| 190 | FailureCode::Invalid, |
| 191 | format!("A run's artifacts are at most {} MB together, and this run's hold {} MB.", RUN_ARTIFACTS_MAX_BYTES / 1_048_576, held / 1_048_576), |
| 192 | )); |
| 193 | } |
| 194 | let existing = self |
| 195 | .artifact_rows("artifacts.run_id = ? AND artifacts.name = ? AND artifacts.status <> 'expired'", &[job.run_id.as_str().into(), a.name.as_str().into()]) |
| 196 | .await?; |
| 197 | if !existing.is_empty() { |
| 198 | if !a.overwrite { |
| 199 | return Ok(fail( |
| 200 | FailureCode::Conflict, |
| 201 | format!("This run already has an artifact named {}. Give the upload another name, or `overwrite: true` to replace it.", a.name), |
| 202 | )); |
| 203 | } |
| 204 | let replaced: Vec<(u64, String)> = existing.iter().map(|row| (row.id as u64, row.object.clone())).collect(); |
| 205 | self.forget_artifacts(&replaced).await?; |
| 206 | } |
| 207 | let now = now_ms(); |
| 208 | let setting = self.retention_setting(&job.repo_id).await?; |
| 209 | let days = retention(a.retention_days, a.expires_at.as_deref().and_then(parse_rfc3339), setting, now); |
| 210 | let expires_at = rfc3339(now + u64::from(days) * DAY_MS); |
| 211 | let format = match a.format.as_deref() { |
| 212 | Some("tgz") => "tgz", |
| 213 | _ => "zip", |
| 214 | }; |
| 215 | // Named by its id once it has one (`object_of` in the API's blobs.rs); |
| 216 | // until then, by a placeholder no other row has. |
| 217 | let placeholder = format!("a/{}/{}", job.repo_id, new_id("art", now)); |
| 218 | let at = rfc3339(now); |
| 219 | #[derive(Deserialize)] |
| 220 | struct Inserted { |
| 221 | id: f64, |
| 222 | } |
| 223 | let inserted = self |
| 224 | .db |
| 225 | .prepare( |
| 226 | "INSERT INTO artifacts (repo_id, namespace, run_id, job_id, name, object, format, size, status, created_at, updated_at, expires_at) |
| 227 | VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'pending', ?9, ?9, ?10) |
| 228 | ON CONFLICT DO NOTHING |
| 229 | RETURNING id", |
| 230 | ) |
| 231 | .bind(&[ |
| 232 | job.repo_id.as_str().into(), |
| 233 | job.namespace.as_str().into(), |
| 234 | job.run_id.as_str().into(), |
| 235 | job.id.as_str().into(), |
| 236 | a.name.as_str().into(), |
| 237 | placeholder.as_str().into(), |
| 238 | format.into(), |
| 239 | (a.size as f64).into(), |
| 240 | at.as_str().into(), |
| 241 | expires_at.as_str().into(), |
| 242 | ])? |
| 243 | .first::<Inserted>(None) |
| 244 | .await?; |
| 245 | let Some(inserted) = inserted else { |
| 246 | return Ok(fail(FailureCode::Conflict, format!("Another job of this run is uploading an artifact named {}.", a.name))); |
| 247 | }; |
| 248 | let id = inserted.id as u64; |
| 249 | let object = object_of(&job.repo_id, id); |
| 250 | self.db.prepare("UPDATE artifacts SET object = ? WHERE id = ?").bind(&[object.as_str().into(), number(id)])?.run().await?; |
| 251 | Ok(Outcome::Ok(ArtifactReservation { id, object, retention_days: days, expires_at })) |
| 252 | } |
| 253 | |
| 254 | /// The pending artifact of a job's run named by id or name. |
| 255 | async fn pending_artifact(&self, run: &str, id: Option<u64>, name: Option<&str>) -> Result<Option<Row>> { |
| 256 | let rows = match (id, name) { |
| 257 | (Some(id), _) => self.artifact_rows("artifacts.id = ? AND artifacts.run_id = ? AND artifacts.status = 'pending'", &[number(id), run.into()]).await?, |
| 258 | (None, Some(name)) => { |
| 259 | self.artifact_rows("artifacts.name = ? AND artifacts.run_id = ? AND artifacts.status = 'pending'", &[name.into(), run.into()]).await? |
| 260 | } |
| 261 | (None, None) => Vec::new(), |
| 262 | }; |
| 263 | Ok(rows.into_iter().next()) |
| 264 | } |
| 265 | |
| 266 | /// `artifact_commit`. |
| 267 | pub async fn artifact_commit(&self, a: ArtifactCommitArgs) -> Result<Outcome<Artifact>> { |
| 268 | let job = check!(self.job_for_credential(&a.job, &a.token).await?); |
| 269 | let Some(row) = self.pending_artifact(&job.run_id, a.id, a.name.as_deref()).await? else { |
| 270 | return Ok(fail(FailureCode::NotFound, "No upload of that artifact is in progress.")); |
| 271 | }; |
| 272 | let id = row.id as u64; |
| 273 | // 0: the size it was measured at as it was stored (the toolkit's |
| 274 | // uploads), rather than one the job says. |
| 275 | let size = if a.size > 0 { a.size } else { row.size as u64 }; |
| 276 | let held = self.run_bytes(&job.run_id, id).await?; |
| 277 | let refused = if size > ARTIFACT_MAX_BYTES { |
| 278 | Some(format!("It is {} MB; an artifact is at most {} MB.", size / 1_048_576, ARTIFACT_MAX_BYTES / 1_048_576)) |
| 279 | } else if held + size > RUN_ARTIFACTS_MAX_BYTES { |
| 280 | Some(format!("A run's artifacts are at most {} MB together, and this run's hold {} MB.", RUN_ARTIFACTS_MAX_BYTES / 1_048_576, held / 1_048_576)) |
| 281 | } else { |
| 282 | None |
| 283 | }; |
| 284 | if let Some(message) = refused { |
| 285 | self.forget_artifacts(&[(id, row.object.clone())]).await?; |
| 286 | return Ok(fail(FailureCode::Invalid, message)); |
| 287 | } |
| 288 | let digest = a.digest.filter(|d| d.starts_with("sha256:") && d.len() == 71); |
| 289 | self.db |
| 290 | .prepare("UPDATE artifacts SET status = 'ready', size = ?, digest = ?, updated_at = ? WHERE id = ? AND status = 'pending'") |
| 291 | .bind(&[(size as f64).into(), crate::optional(digest.as_deref()), rfc3339(now_ms()).into(), number(id)])? |
| 292 | .run() |
| 293 | .await?; |
| 294 | let rows = self.artifact_rows("artifacts.id = ?", &[number(id)]).await?; |
| 295 | Ok(rows.first().map_or_else(|| fail(FailureCode::NotFound, "No such artifact."), |row| Outcome::Ok(row.view()))) |
| 296 | } |
| 297 | |
| 298 | /// `artifact_abort`. |
| 299 | pub async fn artifact_abort(&self, a: JobArtifactsArgs) -> Result<Outcome<bool>> { |
| 300 | let job = check!(self.job_for_credential(&a.job, &a.token).await?); |
| 301 | if let Some(row) = self.pending_artifact(&job.run_id, a.id, a.name.as_deref()).await? { |
| 302 | self.forget_artifacts(&[(row.id as u64, row.object)]).await?; |
| 303 | } |
| 304 | Ok(Outcome::Ok(true)) |
| 305 | } |
| 306 | |
| 307 | /// The ready artifacts of `run`, which must be in the job's repository. |
| 308 | async fn readable_by_job(&self, a: &JobArtifactsArgs) -> Result<Outcome<(crate::plan::JobRow, Vec<Row>)>> { |
| 309 | let job = check!(self.job_for_credential(&a.job, &a.token).await?); |
| 310 | // An id names one artifact of the repository, whichever run made it. |
| 311 | if let (Some(id), None) = (a.id, a.run_id.as_deref().filter(|r| !r.is_empty())) { |
| 312 | let rows = self |
| 313 | .artifact_rows("artifacts.id = ? AND artifacts.repo_id = ? AND artifacts.status = 'ready'", &[number(id), job.repo_id.as_str().into()]) |
| 314 | .await?; |
| 315 | return Ok(Outcome::Ok((job, rows))); |
| 316 | } |
| 317 | let run = a.run_id.clone().filter(|r| !r.is_empty()).unwrap_or_else(|| job.run_id.clone()); |
| 318 | let mut rows = self |
| 319 | .artifact_rows("artifacts.run_id = ? AND artifacts.repo_id = ? AND artifacts.status = 'ready' ORDER BY artifacts.id", &[run.as_str().into(), job.repo_id.as_str().into()]) |
| 320 | .await?; |
| 321 | if rows.is_empty() && run != job.run_id { |
| 322 | let found = self.run_row(&run).await?; |
| 323 | if found.is_none_or(|found| found.repo_id != job.repo_id) { |
| 324 | return Ok(fail(FailureCode::NotFound, format!("This repository has no run {run}."))); |
| 325 | } |
| 326 | } |
| 327 | if let Some(name) = &a.name { |
| 328 | rows.retain(|row| &row.name == name); |
| 329 | } |
| 330 | if let Some(id) = a.id { |
| 331 | rows.retain(|row| row.id as u64 == id); |
| 332 | } |
| 333 | Ok(Outcome::Ok((job, rows))) |
| 334 | } |
| 335 | |
| 336 | /// `job_artifacts`. |
| 337 | pub async fn job_artifacts(&self, a: JobArtifactsArgs) -> Result<Outcome<Vec<Artifact>>> { |
| 338 | let (_, rows) = check!(self.readable_by_job(&a).await?); |
| 339 | Ok(Outcome::Ok(rows.iter().map(Row::view).collect())) |
| 340 | } |
| 341 | |
| 342 | fn blob_of(&self, row: &Row, seconds: u64) -> Outcome<ArtifactBlob> { |
| 343 | match self.download_token("artifact", &(row.id as u64).to_string(), &row.object, seconds) { |
| 344 | Some(blob) => Outcome::Ok(ArtifactBlob { artifact: row.view(), object: row.object.clone(), blob }), |
| 345 | // Without a key there are no download links; the API streams |
| 346 | // the object itself for those who need only its name. |
| 347 | None => Outcome::Ok(ArtifactBlob { artifact: row.view(), object: row.object.clone(), blob: String::new() }), |
| 348 | } |
| 349 | } |
| 350 | |
| 351 | /// `job_artifact`. |
| 352 | pub async fn job_artifact(&self, a: JobArtifactsArgs) -> Result<Outcome<ArtifactBlob>> { |
| 353 | let (_, rows) = check!(self.readable_by_job(&a).await?); |
| 354 | Ok(match rows.last() { |
| 355 | Some(row) => self.blob_of(row, DOWNLOAD_SECONDS), |
| 356 | None => fail(FailureCode::NotFound, "No such artifact, or it has expired."), |
| 357 | }) |
| 358 | } |
| 359 | |
| 360 | /// `job_delete_artifact`. |
| 361 | pub async fn job_delete_artifact(&self, a: JobArtifactsArgs) -> Result<Outcome<Artifact>> { |
| 362 | let job = check!(self.job_for_credential(&a.job, &a.token).await?); |
| 363 | let own = JobArtifactsArgs { run_id: Some(job.run_id.clone()), ..a }; |
| 364 | let (_, rows) = check!(self.readable_by_job(&own).await?); |
| 365 | let Some(row) = rows.last() else { |
| 366 | return Ok(fail(FailureCode::NotFound, "This run has no such artifact.")); |
| 367 | }; |
| 368 | self.forget_artifacts(&[(row.id as u64, row.object.clone())]).await?; |
| 369 | let mut gone = row.view(); |
| 370 | gone.expired = true; |
| 371 | Ok(Outcome::Ok(gone)) |
| 372 | } |
| 373 | |
| 374 | /// `artifacts`. |
| 375 | pub async fn artifacts(&self, a: ArtifactsArgs) -> Result<Outcome<ArtifactList>> { |
| 376 | let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else { |
| 377 | return Ok(fail(FailureCode::NotFound, "There is no such repository.")); |
| 378 | }; |
| 379 | let mut filter = "artifacts.repo_id = ? AND artifacts.status = 'ready'".to_owned(); |
| 380 | let mut values: Vec<JsValue> = vec![repo.id.as_str().into()]; |
| 381 | if let Some(run) = a.run.as_deref().filter(|r| !r.is_empty()) { |
| 382 | if self.run_in(&a.repo, run).await?.into_result().is_err() { |
| 383 | return Ok(fail(FailureCode::NotFound, "No such run.")); |
| 384 | } |
| 385 | filter.push_str(" AND artifacts.run_id = ?"); |
| 386 | values.push(run.into()); |
| 387 | } |
| 388 | if let Some(name) = a.name.as_deref().filter(|n| !n.is_empty()) { |
| 389 | filter.push_str(" AND artifacts.name = ?"); |
| 390 | values.push(name.into()); |
| 391 | } |
| 392 | let rows = self.artifact_rows(&format!("{filter} ORDER BY artifacts.id DESC"), &values).await?; |
| 393 | let per_page = a.per_page.unwrap_or(30).clamp(1, 100) as usize; |
| 394 | let page = a.page.unwrap_or(1).max(1) as usize; |
| 395 | Ok(Outcome::Ok(ArtifactList { |
| 396 | total_count: rows.len() as u64, |
| 397 | artifacts: rows.iter().skip((page - 1) * per_page).take(per_page).map(Row::view).collect(), |
| 398 | })) |
| 399 | } |
| 400 | |
| 401 | /// One ready artifact of a repository the viewer can see. |
| 402 | async fn visible_artifact(&self, a: &ArtifactArgs) -> Result<Outcome<Row>> { |
| 403 | let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else { |
| 404 | return Ok(fail(FailureCode::NotFound, "There is no such repository.")); |
| 405 | }; |
| 406 | let rows = match (a.id, a.run.as_deref(), a.name.as_deref()) { |
| 407 | (Some(id), _, _) => self.artifact_rows("artifacts.id = ? AND artifacts.repo_id = ? AND artifacts.status = 'ready'", &[number(id), repo.id.as_str().into()]).await?, |
| 408 | (None, Some(run), Some(name)) => { |
| 409 | self.artifact_rows( |
| 410 | "artifacts.run_id = ? AND artifacts.name = ? AND artifacts.repo_id = ? AND artifacts.status = 'ready'", |
| 411 | &[run.into(), name.into(), repo.id.as_str().into()], |
| 412 | ) |
| 413 | .await? |
| 414 | } |
| 415 | _ => return Ok(fail(FailureCode::Invalid, "Name the artifact by its id, or by its run and name.")), |
| 416 | }; |
| 417 | Ok(rows.into_iter().next().map_or_else(|| fail(FailureCode::NotFound, "No such artifact, or it has expired."), Outcome::Ok)) |
| 418 | } |
| 419 | |
| 420 | /// `artifact`. |
| 421 | pub async fn artifact(&self, a: ArtifactArgs) -> Result<Outcome<Artifact>> { |
| 422 | let row = check!(self.visible_artifact(&a).await?); |
| 423 | Ok(Outcome::Ok(row.view())) |
| 424 | } |
| 425 | |
| 426 | /// `artifact_download`. |
| 427 | pub async fn artifact_download(&self, a: ArtifactArgs) -> Result<Outcome<ArtifactBlob>> { |
| 428 | let row = check!(self.visible_artifact(&a).await?); |
| 429 | Ok(self.blob_of(&row, LINK_SECONDS)) |
| 430 | } |
| 431 | |
| 432 | /// `delete_artifact`. |
| 433 | pub async fn delete_artifact(&self, a: DeleteArtifactArgs) -> Result<Outcome<Artifact>> { |
| 434 | let repo = check!(self.may(&a.actor, &a.repo, Capability::Run).await?); |
| 435 | let rows = self.artifact_rows("artifacts.id = ? AND artifacts.repo_id = ? AND artifacts.status <> 'expired'", &[number(a.id), repo.id.as_str().into()]).await?; |
| 436 | let Some(row) = rows.into_iter().next() else { |
| 437 | return Ok(fail(FailureCode::NotFound, "No such artifact, or it has expired.")); |
| 438 | }; |
| 439 | self.forget_artifacts(&[(row.id as u64, row.object.clone())]).await?; |
| 440 | let mut gone = row.view(); |
| 441 | gone.expired = true; |
| 442 | Ok(Outcome::Ok(gone)) |
| 443 | } |
| 444 | |
| 445 | /// `artifact_retention`. |
| 446 | pub async fn artifact_retention(&self, a: ArtifactRetentionArgs) -> Result<Outcome<ArtifactRetention>> { |
| 447 | let repo = match a.days { |
| 448 | Some(days) => { |
| 449 | let Some(actor) = a.viewer.clone() else { |
| 450 | return Ok(fail(FailureCode::Unauthenticated, "Sign in to change how long artifacts are kept.")); |
| 451 | }; |
| 452 | let repo = check!(self.may(&actor, &a.repo, Capability::ManageSettings).await?); |
| 453 | if !(1..=ARTIFACT_RETENTION_MAX_DAYS).contains(&days) { |
| 454 | return Ok(fail(FailureCode::Invalid, format!("Artifacts are kept for 1 to {ARTIFACT_RETENTION_MAX_DAYS} days."))); |
| 455 | } |
| 456 | self.db |
| 457 | .prepare( |
| 458 | "INSERT INTO repo_settings (repo_id, artifact_retention_days, updated_at, updated_by) VALUES (?1, ?2, ?3, ?4) |
| 459 | ON CONFLICT (repo_id) DO UPDATE SET artifact_retention_days = ?2, updated_at = ?3, updated_by = ?4", |
| 460 | ) |
| 461 | .bind(&[repo.id.as_str().into(), f64::from(days).into(), rfc3339(now_ms()).into(), actor.id.as_str().into()])? |
| 462 | .run() |
| 463 | .await?; |
| 464 | repo |
| 465 | } |
| 466 | None => match self.visible_repo(&a.repo, &a.viewer).await? { |
| 467 | Some(repo) => repo, |
| 468 | None => return Ok(fail(FailureCode::NotFound, "There is no such repository.")), |
| 469 | }, |
| 470 | }; |
| 471 | Ok(Outcome::Ok(ArtifactRetention { days: self.retention_setting(&repo.id).await?, maximum_allowed_days: ARTIFACT_RETENTION_MAX_DAYS })) |
| 472 | } |
| 473 | |
| 474 | /// Hourly: artifacts past their time, and uploads left unfinished, |
| 475 | /// with their objects. |
| 476 | pub(crate) async fn sweep_artifacts(&self, now: u64) -> Result<()> { |
| 477 | let at = rfc3339(now); |
| 478 | self.db |
| 479 | .prepare("UPDATE artifacts SET status = 'expired', updated_at = ?1 WHERE (status = 'ready' AND expires_at < ?1) OR (status = 'pending' AND created_at < ?2)") |
| 480 | .bind(&[at.as_str().into(), rfc3339(now.saturating_sub(PENDING_MS)).into()])? |
| 481 | .run() |
| 482 | .await?; |
| 483 | #[derive(Deserialize)] |
| 484 | struct Gone { |
| 485 | id: f64, |
| 486 | object: String, |
| 487 | } |
| 488 | let gone = self.db.prepare("SELECT id, object FROM artifacts WHERE status = 'expired' LIMIT 500").all().await?.results::<Gone>()?; |
| 489 | let gone: Vec<(u64, String)> = gone.into_iter().map(|g| (g.id as u64, g.object)).collect(); |
| 490 | self.forget_artifacts(&gone).await |
| 491 | } |
| 492 | } |
| 493 | |
| 494 | #[cfg(test)] |
| 495 | mod tests { |
| 496 | use super::*; |
| 497 | |
| 498 | #[test] |
| 499 | fn names_follow_github_s_rule() { |
| 500 | for good in ["dist", "coverage report", "build (linux, 22)", "my-app_1.0", "résumé"] { |
| 501 | assert!(valid_name(good), "{good}"); |
| 502 | } |
| 503 | for bad in ["", " ", "a/b", "a\\b", "a:b", "a*", "what?", "a\"b", "<x>", "a|b", "a\nb", &"x".repeat(257)] { |
| 504 | assert!(!valid_name(bad), "{bad:?}"); |
| 505 | } |
| 506 | } |
| 507 | |
| 508 | #[test] |
| 509 | fn retention_is_what_was_asked_up_to_the_setting() { |
| 510 | let now = 1_000 * DAY_MS; |
| 511 | // Nothing asked: the repository's setting. |
| 512 | assert_eq!(retention(0, None, 14, now), 14); |
| 513 | assert_eq!(retention(0, None, 90, now), 90); |
| 514 | // Asked for fewer days, or more than the setting allows. |
| 515 | assert_eq!(retention(3, None, 14, now), 3); |
| 516 | assert_eq!(retention(30, None, 14, now), 14); |
| 517 | // The toolkit says when to expire it instead: rounded up to days. |
| 518 | assert_eq!(retention(0, Some(now + 5 * DAY_MS), 30, now), 5); |
| 519 | assert_eq!(retention(0, Some(now + 5 * DAY_MS - 1000), 30, now), 5); |
| 520 | assert_eq!(retention(0, Some(now - 1), 30, now), 30); |
| 521 | // A setting out of range is held to 1 to 90 days. |
| 522 | assert_eq!(retention(0, None, 0, now), 1); |
| 523 | assert_eq!(retention(200, None, 400, now), 90); |
| 524 | } |
| 525 | } |