| 1 | //! Version updates' tables in D1 (migration 0004): when each entry runs, |
| 2 | //! the update pull requests g1t makes, and ignore conditions from comments. |
| 3 | |
| 4 | use g1t_contracts::new_id; |
| 5 | use g1t_contracts::security::{BumpArgs, UpdateState}; |
| 6 | use g1t_contracts::updates::{IgnoreCondition, UpdatePull, UpdatedDependency}; |
| 7 | use g1t_kit::now_ms; |
| 8 | use serde::Deserialize; |
| 9 | use worker::Result; |
| 10 | use worker::wasm_bindgen::JsValue; |
| 11 | |
| 12 | use crate::store::{Store, now, optional}; |
| 13 | |
| 14 | /// Where one `updates` entry's checks stand. |
| 15 | #[derive(Clone, Debug, Deserialize)] |
| 16 | pub struct RunRow { |
| 17 | pub repo_id: String, |
| 18 | pub entry: String, |
| 19 | pub next_run_at: Option<String>, |
| 20 | pub last_checked_at: Option<String>, |
| 21 | pub last_result: Option<String>, |
| 22 | pub last_error: Option<String>, |
| 23 | } |
| 24 | |
| 25 | /// An update pull request, as stored. |
| 26 | #[derive(Clone, Debug, Deserialize)] |
| 27 | pub struct PullRow { |
| 28 | pub id: String, |
| 29 | pub repo_id: String, |
| 30 | pub kind: String, |
| 31 | pub entry: String, |
| 32 | pub ecosystem: String, |
| 33 | pub subject: String, |
| 34 | pub signature: String, |
| 35 | pub group_name: Option<String>, |
| 36 | pub branch: String, |
| 37 | pub title: String, |
| 38 | pub body: String, |
| 39 | pub dependencies: String, |
| 40 | pub bump: String, |
| 41 | pub assignees: String, |
| 42 | pub reviewers: String, |
| 43 | pub state: String, |
| 44 | pub pull: Option<i64>, |
| 45 | pub head: Option<String>, |
| 46 | pub merge_by: Option<String>, |
| 47 | pub error: Option<String>, |
| 48 | pub updated_at: String, |
| 49 | } |
| 50 | |
| 51 | impl PullRow { |
| 52 | pub fn state(&self) -> UpdateState { |
| 53 | UpdateState::parse(&self.state).unwrap_or(UpdateState::Failed) |
| 54 | } |
| 55 | |
| 56 | pub fn pull(&self) -> Option<u32> { |
| 57 | self.pull.map(|n| n.max(0) as u32) |
| 58 | } |
| 59 | |
| 60 | pub fn dependencies(&self) -> Vec<UpdatedDependency> { |
| 61 | serde_json::from_str(&self.dependencies).unwrap_or_default() |
| 62 | } |
| 63 | |
| 64 | pub fn bump(&self) -> Option<BumpArgs> { |
| 65 | serde_json::from_str(&self.bump).ok() |
| 66 | } |
| 67 | |
| 68 | pub fn names(&self, column: &str) -> Vec<String> { |
| 69 | serde_json::from_str(if column == "assignees" { &self.assignees } else { &self.reviewers }).unwrap_or_default() |
| 70 | } |
| 71 | |
| 72 | /// Who asked for it to merge, by username. |
| 73 | pub fn merge_requested_by(&self) -> Option<String> { |
| 74 | self.merge_by.clone() |
| 75 | } |
| 76 | |
| 77 | pub fn to_contract(&self) -> UpdatePull { |
| 78 | UpdatePull { |
| 79 | kind: self.kind.clone(), |
| 80 | entry: self.entry.clone(), |
| 81 | ecosystem: self.ecosystem.clone(), |
| 82 | group: self.group_name.clone(), |
| 83 | branch: self.branch.clone(), |
| 84 | title: self.title.clone(), |
| 85 | state: self.state.clone(), |
| 86 | pull: self.pull(), |
| 87 | dependencies: self.dependencies(), |
| 88 | merge_requested_by: self.merge_requested_by(), |
| 89 | error: self.error.clone(), |
| 90 | updated_at: self.updated_at.clone(), |
| 91 | } |
| 92 | } |
| 93 | } |
| 94 | |
| 95 | /// A pull request to record before its sandbox starts. |
| 96 | pub struct NewPull<'a> { |
| 97 | pub repo_id: &'a str, |
| 98 | pub kind: &'a str, |
| 99 | pub entry: &'a str, |
| 100 | pub ecosystem: &'a str, |
| 101 | pub subject: &'a str, |
| 102 | pub signature: &'a str, |
| 103 | pub group: Option<&'a str>, |
| 104 | pub branch: &'a str, |
| 105 | pub title: &'a str, |
| 106 | pub body: &'a str, |
| 107 | pub dependencies: &'a [UpdatedDependency], |
| 108 | /// Without registries: credentials are never stored. |
| 109 | pub bump: &'a BumpArgs, |
| 110 | pub assignees: &'a [String], |
| 111 | pub reviewers: &'a [String], |
| 112 | } |
| 113 | |
| 114 | #[derive(Deserialize)] |
| 115 | struct IgnoreRow { |
| 116 | ecosystem: String, |
| 117 | dependency: String, |
| 118 | versions: Option<String>, |
| 119 | update_type: Option<String>, |
| 120 | by: String, |
| 121 | pull: Option<i64>, |
| 122 | at: String, |
| 123 | } |
| 124 | |
| 125 | #[derive(Deserialize)] |
| 126 | struct Claimed { |
| 127 | #[allow(dead_code)] |
| 128 | entry: String, |
| 129 | } |
| 130 | |
| 131 | /// What clears a repository's version update rows when it is purged. |
| 132 | pub const PURGED: &[&str] = &[ |
| 133 | "DELETE FROM update_runs WHERE repo_id = ?1", |
| 134 | "DELETE FROM update_pulls WHERE repo_id = ?1", |
| 135 | "DELETE FROM update_ignores WHERE repo_id = ?1", |
| 136 | ]; |
| 137 | |
| 138 | /// A run claimed this long ago and never finished is taken to have died. |
| 139 | const RUN_CLAIM_MS: u64 = 30 * 60 * 1000; |
| 140 | |
| 141 | impl Store { |
| 142 | pub async fn runs(&self, repo_id: &str) -> Result<Vec<RunRow>> { |
| 143 | self.db.prepare("SELECT * FROM update_runs WHERE repo_id = ?").bind(&[repo_id.into()])?.all().await?.results::<RunRow>() |
| 144 | } |
| 145 | |
| 146 | /// Keeps a row for each of `entries` (id, next run) and removes the rest. |
| 147 | pub async fn set_runs(&self, repo_id: &str, entries: &[(String, Option<String>)]) -> Result<()> { |
| 148 | let mut statements = Vec::new(); |
| 149 | for (entry, next) in entries { |
| 150 | statements.push( |
| 151 | self.db |
| 152 | .prepare( |
| 153 | "INSERT INTO update_runs (repo_id, entry, next_run_at) VALUES (?1, ?2, ?3) |
| 154 | ON CONFLICT (repo_id, entry) DO UPDATE SET next_run_at = ?3", |
| 155 | ) |
| 156 | .bind(&[repo_id.into(), entry.as_str().into(), optional(next.as_deref())])?, |
| 157 | ); |
| 158 | } |
| 159 | let keep: Vec<JsValue> = entries.iter().map(|(entry, _)| JsValue::from(entry.as_str())).collect(); |
| 160 | let marks = vec!["?"; keep.len()].join(", "); |
| 161 | let mut binds = vec![JsValue::from(repo_id)]; |
| 162 | binds.extend(keep); |
| 163 | let sql = if entries.is_empty() { |
| 164 | "DELETE FROM update_runs WHERE repo_id = ?".to_owned() |
| 165 | } else { |
| 166 | format!("DELETE FROM update_runs WHERE repo_id = ? AND entry NOT IN ({marks})") |
| 167 | }; |
| 168 | statements.push(self.db.prepare(sql).bind(&binds)?); |
| 169 | self.db.batch(statements).await?; |
| 170 | Ok(()) |
| 171 | } |
| 172 | |
| 173 | /// Entries whose next run has come, oldest first, not being run now. |
| 174 | pub async fn due_runs(&self, at: &str, limit: u32) -> Result<Vec<RunRow>> { |
| 175 | let stale = g1t_contracts::time::rfc3339(now_ms().saturating_sub(RUN_CLAIM_MS)); |
| 176 | self.db |
| 177 | .prepare( |
| 178 | "SELECT * FROM update_runs WHERE next_run_at IS NOT NULL AND next_run_at <= ? |
| 179 | AND (running_at IS NULL OR running_at < ?) ORDER BY next_run_at LIMIT ?", |
| 180 | ) |
| 181 | .bind(&[at.into(), stale.into(), limit.into()])? |
| 182 | .all() |
| 183 | .await? |
| 184 | .results::<RunRow>() |
| 185 | } |
| 186 | |
| 187 | /// Takes an entry's run, unless another sweep already has. |
| 188 | pub async fn claim_run(&self, repo_id: &str, entry: &str) -> Result<bool> { |
| 189 | let stale = g1t_contracts::time::rfc3339(now_ms().saturating_sub(RUN_CLAIM_MS)); |
| 190 | let claimed = self |
| 191 | .db |
| 192 | .prepare( |
| 193 | "INSERT INTO update_runs (repo_id, entry, running_at) VALUES (?1, ?2, ?3) |
| 194 | ON CONFLICT (repo_id, entry) DO UPDATE SET running_at = ?3 |
| 195 | WHERE update_runs.running_at IS NULL OR update_runs.running_at < ?4 |
| 196 | RETURNING entry", |
| 197 | ) |
| 198 | .bind(&[repo_id.into(), entry.into(), now().into(), stale.into()])? |
| 199 | .first::<Claimed>(None) |
| 200 | .await?; |
| 201 | Ok(claimed.is_some()) |
| 202 | } |
| 203 | |
| 204 | pub async fn finish_run(&self, repo_id: &str, entry: &str, next: Option<&str>, result: Option<&str>, error: Option<&str>) -> Result<()> { |
| 205 | self.db |
| 206 | .prepare( |
| 207 | "UPDATE update_runs SET running_at = NULL, last_checked_at = ?, next_run_at = ?, last_result = ?, last_error = ? |
| 208 | WHERE repo_id = ? AND entry = ?", |
| 209 | ) |
| 210 | .bind(&[now().into(), optional(next), optional(result), optional(error), repo_id.into(), entry.into()])? |
| 211 | .run() |
| 212 | .await?; |
| 213 | Ok(()) |
| 214 | } |
| 215 | |
| 216 | // --- Update pull requests ------------------------------------------------- |
| 217 | |
| 218 | pub async fn update_pulls(&self, repo_id: &str) -> Result<Vec<PullRow>> { |
| 219 | self.db |
| 220 | .prepare("SELECT * FROM update_pulls WHERE repo_id = ? ORDER BY requested_at DESC LIMIT 500") |
| 221 | .bind(&[repo_id.into()])? |
| 222 | .all() |
| 223 | .await? |
| 224 | .results::<PullRow>() |
| 225 | } |
| 226 | |
| 227 | /// The latest update pull request made on `branch`. |
| 228 | pub async fn update_pull_by_branch(&self, repo_id: &str, branch: &str) -> Result<Option<PullRow>> { |
| 229 | self.db |
| 230 | .prepare("SELECT * FROM update_pulls WHERE repo_id = ? AND branch = ? ORDER BY requested_at DESC LIMIT 1") |
| 231 | .bind(&[repo_id.into(), branch.into()])? |
| 232 | .first::<PullRow>(None) |
| 233 | .await |
| 234 | } |
| 235 | |
| 236 | pub async fn update_pull_by_number(&self, repo_id: &str, number: u32) -> Result<Option<PullRow>> { |
| 237 | self.db |
| 238 | .prepare("SELECT * FROM update_pulls WHERE repo_id = ? AND pull = ? ORDER BY requested_at DESC LIMIT 1") |
| 239 | .bind(&[repo_id.into(), number.into()])? |
| 240 | .first::<PullRow>(None) |
| 241 | .await |
| 242 | } |
| 243 | |
| 244 | /// Open update pull requests, the least recently looked at first. |
| 245 | pub async fn open_update_pulls(&self, limit: u32) -> Result<Vec<PullRow>> { |
| 246 | self.db |
| 247 | .prepare("SELECT * FROM update_pulls WHERE state = 'open' ORDER BY updated_at LIMIT ?") |
| 248 | .bind(&[limit.into()])? |
| 249 | .all() |
| 250 | .await? |
| 251 | .results::<PullRow>() |
| 252 | } |
| 253 | |
| 254 | /// Update pull requests asked of a sandbox before `before` that never pushed. |
| 255 | pub async fn stalled_update_pulls(&self, before: &str, limit: u32) -> Result<Vec<PullRow>> { |
| 256 | self.db |
| 257 | .prepare("SELECT * FROM update_pulls WHERE state = 'requested' AND updated_at < ? ORDER BY updated_at LIMIT ?") |
| 258 | .bind(&[before.into(), limit.into()])? |
| 259 | .all() |
| 260 | .await? |
| 261 | .results::<PullRow>() |
| 262 | } |
| 263 | |
| 264 | pub async fn add_update_pull(&self, new: &NewPull<'_>) -> Result<String> { |
| 265 | let id = new_id("upd", now_ms()); |
| 266 | let at = now(); |
| 267 | self.db |
| 268 | .prepare( |
| 269 | "INSERT INTO update_pulls (id, repo_id, kind, entry, ecosystem, subject, signature, group_name, branch, title, |
| 270 | body, dependencies, bump, assignees, reviewers, state, requested_at, updated_at) |
| 271 | VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, 'requested', ?16, ?16)", |
| 272 | ) |
| 273 | .bind(&[ |
| 274 | id.as_str().into(), |
| 275 | new.repo_id.into(), |
| 276 | new.kind.into(), |
| 277 | new.entry.into(), |
| 278 | new.ecosystem.into(), |
| 279 | new.subject.into(), |
| 280 | new.signature.into(), |
| 281 | optional(new.group), |
| 282 | new.branch.into(), |
| 283 | new.title.into(), |
| 284 | new.body.into(), |
| 285 | serde_json::to_string(new.dependencies)?.into(), |
| 286 | serde_json::to_string(new.bump)?.into(), |
| 287 | serde_json::to_string(new.assignees)?.into(), |
| 288 | serde_json::to_string(new.reviewers)?.into(), |
| 289 | at.as_str().into(), |
| 290 | ])? |
| 291 | .run() |
| 292 | .await?; |
| 293 | Ok(id) |
| 294 | } |
| 295 | |
| 296 | pub async fn set_update_pull(&self, id: &str, state: UpdateState, pull: Option<u32>, issue: Option<u32>, error: Option<&str>) -> Result<()> { |
| 297 | self.db |
| 298 | .prepare("UPDATE update_pulls SET state = ?, pull = ?, issue = ?, error = ?, updated_at = ? WHERE id = ?") |
| 299 | .bind(&[ |
| 300 | state.as_str().into(), |
| 301 | pull.map_or(JsValue::NULL, JsValue::from), |
| 302 | issue.map_or(JsValue::NULL, JsValue::from), |
| 303 | optional(error), |
| 304 | now().into(), |
| 305 | id.into(), |
| 306 | ])? |
| 307 | .run() |
| 308 | .await?; |
| 309 | Ok(()) |
| 310 | } |
| 311 | |
| 312 | /// Asks for the pull request to be made again on its branch (a rebase |
| 313 | /// or a recreate): it stays open while its sandbox works. |
| 314 | pub async fn remake_update_pull(&self, id: &str, bump: &BumpArgs) -> Result<()> { |
| 315 | self.db |
| 316 | .prepare("UPDATE update_pulls SET bump = ?, head = NULL, error = NULL, updated_at = ? WHERE id = ?") |
| 317 | .bind(&[serde_json::to_string(bump)?.into(), now().into(), id.into()])? |
| 318 | .run() |
| 319 | .await?; |
| 320 | Ok(()) |
| 321 | } |
| 322 | |
| 323 | pub async fn set_update_pull_head(&self, id: &str, head: &str) -> Result<()> { |
| 324 | self.db |
| 325 | .prepare("UPDATE update_pulls SET head = ?, updated_at = ? WHERE id = ?") |
| 326 | .bind(&[head.into(), now().into(), id.into()])? |
| 327 | .run() |
| 328 | .await?; |
| 329 | Ok(()) |
| 330 | } |
| 331 | |
| 332 | /// Marks it looked at, so the sweep moves on to others. |
| 333 | pub async fn touch_update_pull(&self, id: &str) -> Result<()> { |
| 334 | self.db.prepare("UPDATE update_pulls SET updated_at = ? WHERE id = ?").bind(&[now().into(), id.into()])?.run().await?; |
| 335 | Ok(()) |
| 336 | } |
| 337 | |
| 338 | pub async fn set_merge_by(&self, id: &str, by: Option<&str>) -> Result<()> { |
| 339 | self.db.prepare("UPDATE update_pulls SET merge_by = ? WHERE id = ?").bind(&[optional(by), id.into()])?.run().await?; |
| 340 | Ok(()) |
| 341 | } |
| 342 | |
| 343 | // --- Ignore conditions ---------------------------------------------------- |
| 344 | |
| 345 | pub async fn ignores(&self, repo_id: &str) -> Result<Vec<IgnoreCondition>> { |
| 346 | Ok(self |
| 347 | .db |
| 348 | .prepare("SELECT * FROM update_ignores WHERE repo_id = ? ORDER BY dependency, at") |
| 349 | .bind(&[repo_id.into()])? |
| 350 | .all() |
| 351 | .await? |
| 352 | .results::<IgnoreRow>()? |
| 353 | .into_iter() |
| 354 | .map(|row| IgnoreCondition { |
| 355 | ecosystem: row.ecosystem, |
| 356 | dependency: row.dependency, |
| 357 | versions: row.versions, |
| 358 | update_type: row.update_type, |
| 359 | by: row.by, |
| 360 | pull: row.pull.map(|n| n.max(0) as u32), |
| 361 | at: row.at, |
| 362 | }) |
| 363 | .collect()) |
| 364 | } |
| 365 | |
| 366 | pub async fn add_ignore(&self, repo_id: &str, condition: &IgnoreCondition) -> Result<()> { |
| 367 | let key = condition.versions.clone().or_else(|| condition.update_type.clone()).unwrap_or_default(); |
| 368 | self.db |
| 369 | .prepare( |
| 370 | "INSERT INTO update_ignores (repo_id, ecosystem, dependency, condition, versions, update_type, by, pull, at) |
| 371 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) |
| 372 | ON CONFLICT (repo_id, ecosystem, dependency, condition) DO UPDATE SET by = excluded.by, pull = excluded.pull, at = excluded.at", |
| 373 | ) |
| 374 | .bind(&[ |
| 375 | repo_id.into(), |
| 376 | condition.ecosystem.as_str().into(), |
| 377 | condition.dependency.to_lowercase().as_str().into(), |
| 378 | key.as_str().into(), |
| 379 | optional(condition.versions.as_deref()), |
| 380 | optional(condition.update_type.as_deref()), |
| 381 | condition.by.as_str().into(), |
| 382 | condition.pull.map_or(JsValue::NULL, JsValue::from), |
| 383 | now().into(), |
| 384 | ])? |
| 385 | .run() |
| 386 | .await?; |
| 387 | Ok(()) |
| 388 | } |
| 389 | |
| 390 | /// Removes the conditions on `dependency` (`*` for every one) in |
| 391 | /// `ecosystem`, or only those `keep` refuses. Returns how many went. |
| 392 | pub async fn remove_ignores(&self, repo_id: &str, ecosystem: &str, dependency: &str, condition: Option<&str>) -> Result<usize> { |
| 393 | let found = self.ignores(repo_id).await?; |
| 394 | let mut removed = 0; |
| 395 | for item in found { |
| 396 | let named = dependency == "*" || item.dependency.eq_ignore_ascii_case(dependency); |
| 397 | let key = item.versions.clone().or_else(|| item.update_type.clone()).unwrap_or_default(); |
| 398 | if item.ecosystem != ecosystem || !named || condition.is_some_and(|wanted| wanted != key) { |
| 399 | continue; |
| 400 | } |
| 401 | self.db |
| 402 | .prepare("DELETE FROM update_ignores WHERE repo_id = ? AND ecosystem = ? AND dependency = ? AND condition = ?") |
| 403 | .bind(&[repo_id.into(), ecosystem.into(), item.dependency.as_str().into(), key.as_str().into()])? |
| 404 | .run() |
| 405 | .await?; |
| 406 | removed += 1; |
| 407 | } |
| 408 | Ok(removed) |
| 409 | } |
| 410 | } |