Skip to content

g1t/services/security/src/update_store.rs

410 lines15,722 bytesCodeBlame
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
4use g1t_contracts::new_id;
5use g1t_contracts::security::{BumpArgs, UpdateState};
6use g1t_contracts::updates::{IgnoreCondition, UpdatePull, UpdatedDependency};
7use g1t_kit::now_ms;
8use serde::Deserialize;
9use worker::Result;
10use worker::wasm_bindgen::JsValue;
11
12use crate::store::{Store, now, optional};
13
14/// Where one `updates` entry's checks stand.
15#[derive(Clone, Debug, Deserialize)]
16pub 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)]
27pub 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
51impl 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.
96pub 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)]
115struct 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)]
126struct Claimed {
127 #[allow(dead_code)]
128 entry: String,
129}
130
131/// What clears a repository's version update rows when it is purged.
132pub 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.
139const RUN_CLAIM_MS: u64 = 30 * 60 * 1000;
140
141impl 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}