Skip to content

g1t/services/security/src/store.rs

1,154 lines42,661 bytesCodeBlame
1//! The security service's tables in D1.
2
3use std::collections::HashMap;
4
5use g1t_contracts::security::{
6 AlertActivity, DismissReason, NewSecret, ScanState, SecretCounts, SecretFinding, SecretStatus, SecurityUpdate,
7 SeverityCounts, UpdateState, VersionUpdatesState, VulnStatus, Vulnerability,
8};
9use g1t_contracts::time::rfc3339;
10use g1t_contracts::new_id;
11use g1t_kit::now_ms;
12use serde::Deserialize;
13use worker::wasm_bindgen::JsValue;
14use worker::{D1Database, Result};
15
16pub fn optional(value: Option<&str>) -> JsValue {
17 value.map_or(JsValue::NULL, JsValue::from)
18}
19
20pub fn now() -> String {
21 rfc3339(now_ms())
22}
23
24/// A repository purged (`repo.purged`): every row kept for it, `?1` its
25/// id. Advisories are shared by every repository, so they stay.
26pub const PURGED: &[&str] = &[
27 "DELETE FROM secrets WHERE repo_id = ?1",
28 "DELETE FROM vulnerabilities WHERE repo_id = ?1",
29 "DELETE FROM upgrades WHERE repo_id = ?1",
30 "DELETE FROM updates WHERE repo_id = ?1",
31 "DELETE FROM alert_activity WHERE repo_id = ?1",
32 "DELETE FROM push_scans WHERE repo_id = ?1",
33 "DELETE FROM repos WHERE repo_id = ?1",
34];
35
36#[derive(Clone, Deserialize)]
37pub struct RepoRow {
38 pub repo_id: String,
39 pub namespace: String,
40 pub name: String,
41 pub upkeep: i64,
42 pub history: String,
43 pub history_cursor: Option<String>,
44 pub history_commits: i64,
45 pub history_finished_at: Option<String>,
46 pub deps_scanned_at: Option<String>,
47 pub deps_error: Option<String>,
48 pub lockfiles: String,
49 #[serde(default)]
50 pub version_updates: Option<String>,
51}
52
53impl RepoRow {
54 pub fn scan_state(&self) -> ScanState {
55 ScanState {
56 history: self.history.clone(),
57 commits_scanned: self.history_commits.max(0) as u32,
58 history_finished_at: self.history_finished_at.clone(),
59 dependencies_scanned_at: self.deps_scanned_at.clone(),
60 dependencies_error: self.deps_error.clone(),
61 lockfiles: serde_json::from_str(&self.lockfiles).unwrap_or_default(),
62 }
63 }
64
65 /// What `.g1t/dependencies.yml` said when it was last read.
66 pub fn version_updates(&self) -> VersionUpdatesState {
67 self.version_updates
68 .as_deref()
69 .and_then(|json| serde_json::from_str(json).ok())
70 .unwrap_or_default()
71 }
72}
73
74#[derive(Deserialize)]
75struct SecretRow {
76 id: String,
77 repo_id: String,
78 kind: String,
79 path: String,
80 line: i64,
81 commit_hash: String,
82 preview: String,
83 status: String,
84 source: String,
85 found_by: Option<String>,
86 found_at: String,
87 decided_by: Option<String>,
88 reason: Option<String>,
89 decided_at: Option<String>,
90 dismiss_reason: Option<String>,
91 test_value: Option<String>,
92 // The security suite's (migration 0004).
93 #[serde(default)]
94 validity: Option<String>,
95 #[serde(default)]
96 validity_checked_at: Option<String>,
97 #[serde(default)]
98 bypass_reason: Option<String>,
99 #[serde(default)]
100 bypass_comment: Option<String>,
101 #[serde(default)]
102 bypassed_by: Option<String>,
103 #[serde(default)]
104 bypassed_at: Option<String>,
105 #[serde(default)]
106 bypass_approved_by: Option<String>,
107 #[serde(default)]
108 pattern_id: Option<String>,
109 #[serde(default)]
110 pattern_name: Option<String>,
111 #[serde(default)]
112 location_count: Option<i64>,
113}
114
115/// Secrets with how many places each was found.
116const SECRET_SELECT: &str =
117 "SELECT s.*, (SELECT count(*) FROM secret_locations l WHERE l.secret_id = s.id) AS location_count FROM secrets s";
118
119impl From<SecretRow> for SecretFinding {
120 fn from(row: SecretRow) -> Self {
121 let status = SecretStatus::parse(&row.status).unwrap_or(SecretStatus::Open);
122 let bypass = match (row.bypass_reason.as_deref().and_then(g1t_contracts::security_suite::BypassReason::parse), row.bypassed_by, row.bypassed_at) {
123 (Some(reason), Some(by), Some(at)) => Some(g1t_contracts::security_suite::Bypass {
124 reason,
125 comment: row.bypass_comment,
126 by,
127 at,
128 approved_by: row.bypass_approved_by,
129 }),
130 _ => None,
131 };
132 SecretFinding {
133 label: crate::secret_alerts::label_of(&row.kind, row.pattern_name.as_deref()),
134 validity: row.validity,
135 validity_checked_at: row.validity_checked_at,
136 bypass,
137 pattern_id: row.pattern_id,
138 pattern_name: row.pattern_name,
139 locations: row.location_count.unwrap_or(0).max(0) as u32,
140 id: row.id,
141 repo_id: row.repo_id,
142 kind: row.kind,
143 path: row.path,
144 line: row.line.max(0) as u32,
145 commit: row.commit_hash,
146 preview: row.preview,
147 status,
148 source: row.source,
149 found_by: row.found_by,
150 found_at: row.found_at,
151 decided_by: row.decided_by,
152 reason: row.reason,
153 decided_at: row.decided_at,
154 dismissed_reason: row.dismiss_reason.as_deref().and_then(DismissReason::parse),
155 test_value: row.test_value,
156 state: status.state(),
157 }
158 }
159}
160
161#[derive(Deserialize)]
162pub struct VulnRow {
163 pub id: String,
164 pub repo_id: String,
165 pub ecosystem: String,
166 pub package: String,
167 pub version: String,
168 pub manifest: String,
169 pub osv_id: String,
170 pub advisory: String,
171 pub summary: String,
172 pub severity: String,
173 pub fixed_version: Option<String>,
174 pub status: String,
175 pub found_at: String,
176 pub fixed_at: Option<String>,
177 pub number: Option<i64>,
178 #[serde(default)]
179 pub dismiss_reason: Option<String>,
180 #[serde(default)]
181 pub dismiss_comment: Option<String>,
182 #[serde(default)]
183 pub dismissed_by: Option<String>,
184 #[serde(default)]
185 pub dismissed_at: Option<String>,
186}
187
188impl From<VulnRow> for Vulnerability {
189 fn from(row: VulnRow) -> Self {
190 let status = VulnStatus::parse(&row.status).unwrap_or(VulnStatus::Open);
191 let dismissed = status == VulnStatus::Dismissed;
192 Vulnerability {
193 id: row.id,
194 repo_id: row.repo_id,
195 ecosystem: row.ecosystem,
196 package: row.package,
197 version: row.version,
198 manifest: row.manifest,
199 advisory: row.advisory,
200 osv_id: row.osv_id,
201 summary: row.summary,
202 severity: row.severity,
203 fixed_version: row.fixed_version,
204 status,
205 issue: row.number.map(|number| number as u32),
206 found_at: row.found_at,
207 fixed_at: row.fixed_at,
208 state: status.state(),
209 dismissed_by: row.dismissed_by.filter(|_| dismissed),
210 dismissed_reason: row.dismiss_reason.as_deref().and_then(DismissReason::parse).filter(|_| dismissed),
211 dismissed_comment: row.dismiss_comment.filter(|_| dismissed),
212 dismissed_at: row.dismissed_at.filter(|_| dismissed),
213 update: None,
214 }
215 }
216}
217
218/// A security update, as stored.
219#[derive(Clone, Deserialize)]
220pub struct UpdateRow {
221 pub repo_id: String,
222 pub ecosystem: String,
223 pub package: String,
224 pub target: String,
225 pub state: String,
226 pub branch: Option<String>,
227 pub pull: Option<i64>,
228 pub issue: Option<i64>,
229 pub error: Option<String>,
230 pub updated_at: String,
231}
232
233impl UpdateRow {
234 pub fn state(&self) -> UpdateState {
235 UpdateState::parse(&self.state).unwrap_or(UpdateState::Failed)
236 }
237
238 pub fn pull(&self) -> Option<u32> {
239 self.pull.map(|n| n.max(0) as u32)
240 }
241
242 pub fn issue(&self) -> Option<u32> {
243 self.issue.map(|n| n.max(0) as u32)
244 }
245
246 pub fn to_contract(&self) -> SecurityUpdate {
247 SecurityUpdate {
248 state: self.state(),
249 target: self.target.clone(),
250 branch: self.branch.clone(),
251 pull: self.pull(),
252 issue: self.issue(),
253 error: self.error.clone(),
254 updated_at: self.updated_at.clone(),
255 }
256 }
257}
258
259#[derive(Deserialize)]
260struct ActivityRow {
261 id: String,
262 alert_id: String,
263 action: String,
264 actor: Option<String>,
265 reason: Option<String>,
266 comment: Option<String>,
267 number: Option<i64>,
268 at: String,
269}
270
271impl From<ActivityRow> for AlertActivity {
272 fn from(row: ActivityRow) -> Self {
273 AlertActivity {
274 id: row.id,
275 alert_id: row.alert_id,
276 action: row.action,
277 actor: row.actor,
278 reason: row.reason.as_deref().and_then(DismissReason::parse),
279 comment: row.comment,
280 number: row.number.map(|n| n.max(0) as u32),
281 at: row.at,
282 }
283 }
284}
285
286/// A push too large to scan before it was stored, scanned after it landed.
287#[derive(Clone, Deserialize)]
288pub struct PushScanRow {
289 pub id: String,
290 pub repo_id: String,
291 pub git_ref: String,
292 pub head: String,
293 pub base: Option<String>,
294 pub cursor: Option<String>,
295 pub pusher: Option<String>,
296 pub pages: i64,
297}
298
299/// One action on an alert, to record.
300pub struct Activity<'a> {
301 pub alert_id: &'a str,
302 pub action: &'a str,
303 pub actor: Option<&'a str>,
304 pub reason: Option<DismissReason>,
305 pub comment: Option<&'a str>,
306 pub number: Option<u32>,
307}
308
309#[derive(Deserialize)]
310pub struct UpgradeRow {
311 pub number: i64,
312}
313
314#[derive(Deserialize)]
315struct CountRow {
316 severity: String,
317 n: i64,
318}
319
320#[derive(Deserialize)]
321struct SecretCountRow {
322 status: String,
323 test: i64,
324 n: i64,
325}
326
327#[derive(Deserialize)]
328struct NumberRow {
329 n: i64,
330}
331
332#[derive(Deserialize)]
333struct FingerprintRow {
334 fingerprint: String,
335 id: String,
336 status: String,
337}
338
339#[derive(Deserialize)]
340struct BodyRow {
341 body: String,
342 fetched_at: String,
343}
344
345pub struct Store {
346 pub db: D1Database,
347}
348
349const VULN_COLUMNS: &str = "v.id, v.repo_id, v.ecosystem, v.package, v.version, v.manifest, v.osv_id, v.advisory,
350 v.summary, v.severity, v.fixed_version, v.status, v.found_at, v.fixed_at, v.dismiss_reason, v.dismiss_comment,
351 v.dismissed_by, v.dismissed_at, COALESCE(n.issue, u.number) AS number";
352
353/// The issue an upgrade has: the security update's, when it went to
354/// g1t, or one opened before security updates.
355const VULN_JOINS: &str = "LEFT JOIN upgrades u ON u.repo_id = v.repo_id AND u.ecosystem = v.ecosystem AND u.package = v.package
356 LEFT JOIN updates n ON n.repo_id = v.repo_id AND n.ecosystem = v.ecosystem AND n.package = v.package";
357
358impl Store {
359 pub async fn repo(&self, repo_id: &str) -> Result<Option<RepoRow>> {
360 self.db
361 .prepare("SELECT * FROM repos WHERE repo_id = ?")
362 .bind(&[repo_id.into()])?
363 .first::<RepoRow>(None)
364 .await
365 }
366
367 /// Records a repository the first time it is seen, and keeps its
368 /// address current. Returns its row.
369 pub async fn register(&self, repo_id: &str, namespace: &str, name: &str) -> Result<RepoRow> {
370 self.db
371 .prepare(
372 "INSERT INTO repos (repo_id, namespace, name, created_at) VALUES (?1, ?2, ?3, ?4)
373 ON CONFLICT (repo_id) DO UPDATE SET namespace = ?2, name = ?3",
374 )
375 .bind(&[repo_id.into(), namespace.into(), name.into(), now().into()])?
376 .run()
377 .await?;
378 self.repo(repo_id)
379 .await?
380 .ok_or_else(|| worker::Error::RustError("the repository was not recorded".into()))
381 }
382
383 /// A repository's path changed: transferred or renamed.
384 pub async fn moved(&self, repo_id: &str, namespace: &str, name: &str) -> Result<()> {
385 self.db
386 .prepare("UPDATE repos SET namespace = ?, name = ? WHERE repo_id = ?")
387 .bind(&[namespace.into(), name.into(), repo_id.into()])?
388 .run()
389 .await?;
390 Ok(())
391 }
392
393 pub async fn rename_namespace(&self, stale: &[String], current: &str) -> Result<()> {
394 for slug in stale {
395 self.db
396 .prepare("UPDATE repos SET namespace = ? WHERE namespace = ?")
397 .bind(&[current.into(), slug.as_str().into()])?
398 .run()
399 .await?;
400 }
401 Ok(())
402 }
403
404 pub async fn in_namespace(&self, namespace: &str) -> Result<Vec<RepoRow>> {
405 self.db
406 .prepare("SELECT * FROM repos WHERE namespace = ? ORDER BY name")
407 .bind(&[namespace.into()])?
408 .all()
409 .await?
410 .results::<RepoRow>()
411 }
412
413 /// A repository purged: everything found in it goes. `?1` its id.
414 pub async fn purge(&self, repo_id: &str) -> Result<()> {
415 let mut batch = Vec::with_capacity(PURGED.len());
416 for sql in PURGED.iter().chain(crate::update_store::PURGED) {
417 batch.push(self.db.prepare(*sql).bind(&[repo_id.into()])?);
418 }
419 self.db.batch(batch).await?;
420 Ok(())
421 }
422
423 /// Records that a repository's daily dependency read was skipped, and
424 /// why, so the sweep moves on to others until the next day.
425 pub async fn skip_dependencies(&self, repo_id: &str, why: &str) -> Result<()> {
426 self.db
427 .prepare("UPDATE repos SET deps_scanned_at = ?, deps_error = ? WHERE repo_id = ?")
428 .bind(&[now().into(), why.into(), repo_id.into()])?
429 .run()
430 .await?;
431 Ok(())
432 }
433
434 /// Repositories whose history still has to be scanned, oldest first.
435 pub async fn unfinished_histories(&self, limit: u32) -> Result<Vec<RepoRow>> {
436 self.db
437 .prepare("SELECT * FROM repos WHERE history IN ('pending', 'running', 'stopped') ORDER BY created_at LIMIT ?")
438 .bind(&[limit.into()])?
439 .all()
440 .await?
441 .results::<RepoRow>()
442 }
443
444 /// Repositories whose dependencies were last read before `before`.
445 pub async fn stale_dependencies(&self, before: &str, limit: u32) -> Result<Vec<RepoRow>> {
446 self.db
447 .prepare(
448 "SELECT * FROM repos WHERE deps_scanned_at IS NULL OR deps_scanned_at < ?
449 ORDER BY deps_scanned_at LIMIT ?",
450 )
451 .bind(&[before.into(), limit.into()])?
452 .all()
453 .await?
454 .results::<RepoRow>()
455 }
456
457 pub async fn set_history(&self, repo_id: &str, state: &str, cursor: Option<&str>, commits: u32) -> Result<()> {
458 let finished = (state == "done").then(now);
459 self.db
460 .prepare(
461 "UPDATE repos SET history = ?, history_cursor = ?, history_commits = history_commits + ?,
462 history_finished_at = COALESCE(?, history_finished_at)
463 WHERE repo_id = ?",
464 )
465 .bind(&[state.into(), optional(cursor), commits.into(), optional(finished.as_deref()), repo_id.into()])?
466 .run()
467 .await?;
468 Ok(())
469 }
470
471 pub async fn restart_history(&self, repo_id: &str) -> Result<()> {
472 self.db
473 .prepare(
474 "UPDATE repos SET history = 'pending', history_cursor = NULL, history_commits = 0,
475 history_finished_at = NULL WHERE repo_id = ?",
476 )
477 .bind(&[repo_id.into()])?
478 .run()
479 .await?;
480 Ok(())
481 }
482
483 pub async fn set_dependencies_scanned(
484 &self,
485 repo_id: &str,
486 commit: Option<&str>,
487 lockfiles: &[String],
488 error: Option<&str>,
489 ) -> Result<()> {
490 self.db
491 .prepare("UPDATE repos SET deps_scanned_at = ?, deps_commit = ?, lockfiles = ?, deps_error = ? WHERE repo_id = ?")
492 .bind(&[now().into(), optional(commit), serde_json::to_string(lockfiles)?.into(), optional(error), repo_id.into()])?
493 .run()
494 .await?;
495 Ok(())
496 }
497
498 pub async fn set_version_updates(&self, repo_id: &str, state: &VersionUpdatesState) -> Result<()> {
499 self.db
500 .prepare("UPDATE repos SET version_updates = ? WHERE repo_id = ?")
501 .bind(&[serde_json::to_string(state)?.into(), repo_id.into()])?
502 .run()
503 .await?;
504 Ok(())
505 }
506
507 pub async fn set_upkeep(&self, repo_id: &str, enabled: bool, by: &str) -> Result<()> {
508 self.db
509 .prepare("UPDATE repos SET upkeep = ?, upkeep_by = ?, upkeep_at = ? WHERE repo_id = ?")
510 .bind(&[i32::from(enabled).into(), by.into(), now().into(), repo_id.into()])?
511 .run()
512 .await?;
513 Ok(())
514 }
515
516 // --- Secrets ------------------------------------------------------------
517
518 pub async fn secrets(&self, repo_id: &str) -> Result<Vec<SecretFinding>> {
519 let rows = self
520 .db
521 .prepare(format!(
522 "{SECRET_SELECT} WHERE s.repo_id = ?
523 ORDER BY CASE status WHEN 'open' THEN 0 WHEN 'blocked' THEN 1 WHEN 'allowed' THEN 2 ELSE 3 END,
524 test_value IS NOT NULL, found_at DESC LIMIT 500"
525 ))
526 .bind(&[repo_id.into()])?
527 .all()
528 .await?
529 .results::<SecretRow>()?;
530 Ok(rows.into_iter().map(SecretFinding::from).collect())
531 }
532
533 pub async fn secret(&self, repo_id: &str, id: &str) -> Result<Option<SecretFinding>> {
534 Ok(self
535 .db
536 .prepare(format!("{SECRET_SELECT} WHERE s.repo_id = ? AND s.id = ?"))
537 .bind(&[repo_id.into(), id.into()])?
538 .first::<SecretRow>(None)
539 .await?
540 .map(SecretFinding::from))
541 }
542
543 /// The findings already known for these fingerprints: fingerprint, id
544 /// and status.
545 pub async fn known(&self, repo_id: &str, fingerprints: &[String]) -> Result<Vec<(String, String, String)>> {
546 let mut known = Vec::new();
547 for chunk in fingerprints.chunks(90) {
548 let marks = vec!["?"; chunk.len()].join(", ");
549 let mut binds: Vec<JsValue> = vec![repo_id.into()];
550 binds.extend(chunk.iter().map(|fingerprint| JsValue::from(fingerprint.as_str())));
551 let rows = self
552 .db
553 .prepare(format!("SELECT fingerprint, id, status FROM secrets WHERE repo_id = ? AND fingerprint IN ({marks})"))
554 .bind(&binds)?
555 .all()
556 .await?
557 .results::<FingerprintRow>()?;
558 known.extend(rows.into_iter().map(|row| (row.fingerprint, row.id, row.status)));
559 }
560 Ok(known)
561 }
562
563 /// Records secrets not seen before in this repository, with `status`.
564 /// A secret already known keeps its record and its decision.
565 pub async fn add_secrets(
566 &self,
567 repo_id: &str,
568 secrets: &[NewSecret],
569 status: SecretStatus,
570 source: &str,
571 found_by: Option<&str>,
572 ) -> Result<()> {
573 if secrets.is_empty() {
574 return Ok(());
575 }
576 let now_ms = now_ms();
577 let found_at = rfc3339(now_ms);
578 let statements = secrets
579 .iter()
580 .map(|secret| {
581 self.db
582 .prepare(
583 "INSERT INTO secrets (id, repo_id, fingerprint, kind, path, line, commit_hash, preview,
584 status, source, found_by, found_at, test_value)
585 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
586 ON CONFLICT (repo_id, fingerprint) DO UPDATE SET test_value = excluded.test_value",
587 )
588 .bind(&[
589 new_id("sec", now_ms).into(),
590 repo_id.into(),
591 secret.fingerprint.as_str().into(),
592 secret.kind.as_str().into(),
593 secret.path.as_str().into(),
594 secret.line.into(),
595 secret.commit.as_str().into(),
596 secret.preview.as_str().into(),
597 status.as_str().into(),
598 source.into(),
599 optional(found_by),
600 found_at.as_str().into(),
601 optional(secret.test_value.as_deref()),
602 ])
603 })
604 .collect::<Result<Vec<_>>>()?;
605 self.db.batch(statements).await?;
606 Ok(())
607 }
608
609 /// A secret found in history that was only ever blocked has landed
610 /// after all (it was allowed, then pushed): it is open now unless
611 /// someone allowed it.
612 pub async fn landed(&self, repo_id: &str, fingerprints: &[String]) -> Result<()> {
613 for chunk in fingerprints.chunks(90) {
614 let marks = vec!["?"; chunk.len()].join(", ");
615 let mut binds: Vec<JsValue> = vec![repo_id.into()];
616 binds.extend(chunk.iter().map(|fingerprint| JsValue::from(fingerprint.as_str())));
617 self.db
618 .prepare(format!(
619 "UPDATE secrets SET status = 'open', source = 'history'
620 WHERE repo_id = ? AND status = 'blocked' AND fingerprint IN ({marks})"
621 ))
622 .bind(&binds)?
623 .run()
624 .await?;
625 }
626 Ok(())
627 }
628
629 /// Dismisses a secret: allowed, or resolved when it was revoked.
630 pub async fn dismiss_secret(&self, repo_id: &str, id: &str, reason: DismissReason, by: &str, comment: Option<&str>) -> Result<()> {
631 self.db
632 .prepare(
633 "UPDATE secrets SET status = ?, decided_by = ?, reason = ?, decided_at = ?, dismiss_reason = ?
634 WHERE repo_id = ? AND id = ?",
635 )
636 .bind(&[
637 reason.secret_status().as_str().into(),
638 by.into(),
639 optional(comment),
640 now().into(),
641 reason.as_str().into(),
642 repo_id.into(),
643 id.into(),
644 ])?
645 .run()
646 .await?;
647 Ok(())
648 }
649
650 /// Opens a secret again, as `status` (open, or blocked for one that
651 /// never landed).
652 pub async fn reopen_secret(&self, repo_id: &str, id: &str, status: SecretStatus) -> Result<()> {
653 self.db
654 .prepare(
655 "UPDATE secrets SET status = ?, decided_by = NULL, reason = NULL, decided_at = NULL, dismiss_reason = NULL
656 WHERE repo_id = ? AND id = ?",
657 )
658 .bind(&[status.as_str().into(), repo_id.into(), id.into()])?
659 .run()
660 .await?;
661 Ok(())
662 }
663
664 // --- Vulnerabilities ----------------------------------------------------
665
666 pub async fn vulnerabilities(&self, repo_id: &str) -> Result<Vec<Vulnerability>> {
667 let rows = self
668 .db
669 .prepare(format!(
670 "SELECT {VULN_COLUMNS} FROM vulnerabilities v
671 {VULN_JOINS}
672 WHERE v.repo_id = ?
673 ORDER BY CASE v.status WHEN 'open' THEN 0 WHEN 'dismissed' THEN 1 ELSE 2 END,
674 CASE v.severity WHEN 'critical' THEN 0 WHEN 'high' THEN 1 WHEN 'medium' THEN 2 WHEN 'low' THEN 3 ELSE 4 END,
675 v.package LIMIT 1000"
676 ))
677 .bind(&[repo_id.into()])?
678 .all()
679 .await?
680 .results::<VulnRow>()?;
681 let updates: HashMap<(String, String), UpdateRow> = self
682 .updates(repo_id)
683 .await?
684 .into_iter()
685 .map(|row| ((row.ecosystem.clone(), row.package.clone()), row))
686 .collect();
687 Ok(rows
688 .into_iter()
689 .map(|row| {
690 let update = updates.get(&(row.ecosystem.clone(), row.package.clone())).map(UpdateRow::to_contract);
691 Vulnerability { update, ..Vulnerability::from(row) }
692 })
693 .collect())
694 }
695
696 pub async fn vulnerability(&self, repo_id: &str, id: &str) -> Result<Option<Vulnerability>> {
697 let row = self
698 .db
699 .prepare(format!("SELECT {VULN_COLUMNS} FROM vulnerabilities v {VULN_JOINS} WHERE v.repo_id = ? AND v.id = ?"))
700 .bind(&[repo_id.into(), id.into()])?
701 .first::<VulnRow>(None)
702 .await?;
703 let Some(row) = row else { return Ok(None) };
704 let update = self.update(repo_id, &row.ecosystem, &row.package).await?.map(|row| row.to_contract());
705 Ok(Some(Vulnerability { update, ..Vulnerability::from(row) }))
706 }
707
708 pub async fn dismiss_vulnerability(
709 &self,
710 repo_id: &str,
711 id: &str,
712 reason: DismissReason,
713 by: &str,
714 comment: Option<&str>,
715 ) -> Result<()> {
716 self.db
717 .prepare(
718 "UPDATE vulnerabilities SET status = CASE status WHEN 'fixed' THEN 'fixed' ELSE 'dismissed' END,
719 dismiss_reason = ?, dismiss_comment = ?, dismissed_by = ?, dismissed_at = ?
720 WHERE repo_id = ? AND id = ?",
721 )
722 .bind(&[reason.as_str().into(), optional(comment), by.into(), now().into(), repo_id.into(), id.into()])?
723 .run()
724 .await?;
725 Ok(())
726 }
727
728 pub async fn reopen_vulnerability(&self, repo_id: &str, id: &str) -> Result<()> {
729 self.db
730 .prepare(
731 "UPDATE vulnerabilities SET status = CASE status WHEN 'dismissed' THEN 'open' ELSE status END,
732 dismiss_reason = NULL, dismiss_comment = NULL, dismissed_by = NULL, dismissed_at = NULL
733 WHERE repo_id = ? AND id = ?",
734 )
735 .bind(&[repo_id.into(), id.into()])?
736 .run()
737 .await?;
738 Ok(())
739 }
740
741 /// The ids of a package's open vulnerabilities.
742 pub async fn open_ids(&self, repo_id: &str, ecosystem: &str, package: &str) -> Result<Vec<String>> {
743 #[derive(Deserialize)]
744 struct IdRow {
745 id: String,
746 }
747 Ok(self
748 .db
749 .prepare("SELECT id FROM vulnerabilities WHERE repo_id = ? AND ecosystem = ? AND package = ? AND status = 'open'")
750 .bind(&[repo_id.into(), ecosystem.into(), package.into()])?
751 .all()
752 .await?
753 .results::<IdRow>()?
754 .into_iter()
755 .map(|row| row.id)
756 .collect())
757 }
758
759 pub async fn open_vulnerabilities(&self, repo_id: &str) -> Result<Vec<VulnRow>> {
760 self.db
761 .prepare(format!(
762 "SELECT {VULN_COLUMNS} FROM vulnerabilities v
763 {VULN_JOINS}
764 WHERE v.repo_id = ? AND v.status = 'open'"
765 ))
766 .bind(&[repo_id.into()])?
767 .all()
768 .await?
769 .results::<VulnRow>()
770 }
771
772 /// Replaces what is known about a repository's dependencies with what a
773 /// scan found: new findings open, findings no longer true fixed, and
774 /// findings that came back open again, or dismissed again when someone
775 /// had dismissed them.
776 pub async fn replace_vulnerabilities(&self, repo_id: &str, found: &[Vulnerability]) -> Result<()> {
777 let now_ms = now_ms();
778 let now = rfc3339(now_ms);
779 let mut statements = vec![self
780 .db
781 .prepare(
782 "UPDATE vulnerabilities SET status = 'fixed', fixed_at = ? WHERE repo_id = ? AND status IN ('open', 'dismissed')",
783 )
784 .bind(&[now.as_str().into(), repo_id.into()])?];
785 for vuln in found {
786 statements.push(
787 self.db
788 .prepare(
789 "INSERT INTO vulnerabilities (id, repo_id, ecosystem, package, version, manifest, osv_id, advisory,
790 summary, severity, fixed_version, status, found_at)
791 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'open', ?)
792 ON CONFLICT (repo_id, ecosystem, package, version, manifest, osv_id) DO UPDATE SET
793 status = CASE WHEN vulnerabilities.dismiss_reason IS NULL THEN 'open' ELSE 'dismissed' END,
794 fixed_at = NULL, advisory = excluded.advisory, summary = excluded.summary,
795 severity = excluded.severity, fixed_version = excluded.fixed_version",
796 )
797 .bind(&[
798 new_id("vul", now_ms).into(),
799 repo_id.into(),
800 vuln.ecosystem.as_str().into(),
801 vuln.package.as_str().into(),
802 vuln.version.as_str().into(),
803 vuln.manifest.as_str().into(),
804 vuln.osv_id.as_str().into(),
805 vuln.advisory.as_str().into(),
806 vuln.summary.as_str().into(),
807 vuln.severity.as_str().into(),
808 optional(vuln.fixed_version.as_deref()),
809 now.as_str().into(),
810 ])?,
811 );
812 }
813 // D1 runs a batch as one transaction, so readers never see the
814 // moment between marking everything fixed and opening it again.
815 self.db.batch(statements).await?;
816 Ok(())
817 }
818
819 pub async fn counts(&self, repo_id: &str) -> Result<(SeverityCounts, u32, u32)> {
820 let rows = self
821 .db
822 .prepare("SELECT severity, count(*) AS n FROM vulnerabilities WHERE repo_id = ? AND status = 'open' GROUP BY severity")
823 .bind(&[repo_id.into()])?
824 .all()
825 .await?
826 .results::<CountRow>()?;
827 let mut counts = SeverityCounts::default();
828 let mut vulns = 0;
829 for row in rows {
830 let n = row.n.max(0) as u32;
831 vulns += n;
832 match row.severity.as_str() {
833 "critical" => counts.critical += n,
834 "high" => counts.high += n,
835 "medium" => counts.medium += n,
836 "low" => counts.low += n,
837 _ => counts.unknown += n,
838 }
839 }
840 let secrets = self.secret_counts(repo_id).await?;
841 // A secret in the history that looks real is the worst thing a
842 // repository can hold. One stopped at a push never landed, and a
843 // likely test value is no danger: neither is counted.
844 counts.critical += secrets.open;
845 Ok((counts, secrets.open + secrets.blocked, vulns))
846 }
847
848 pub async fn secret_counts(&self, repo_id: &str) -> Result<SecretCounts> {
849 let rows = self
850 .db
851 .prepare(
852 "SELECT status, (test_value IS NOT NULL) AS test, count(*) AS n FROM secrets WHERE repo_id = ?
853 GROUP BY status, test",
854 )
855 .bind(&[repo_id.into()])?
856 .all()
857 .await?
858 .results::<SecretCountRow>()?;
859 Ok(secret_counts(rows.iter().map(|row| (row.status.as_str(), row.test != 0, row.n.max(0) as u32))))
860 }
861
862 // --- Activity -----------------------------------------------------------
863
864 pub async fn record(&self, repo_id: &str, activity: &[Activity<'_>]) -> Result<()> {
865 if activity.is_empty() {
866 return Ok(());
867 }
868 let now_ms = now_ms();
869 let at = rfc3339(now_ms);
870 let statements = activity
871 .iter()
872 .map(|item| {
873 self.db
874 .prepare(
875 "INSERT INTO alert_activity (id, repo_id, alert_id, action, actor, reason, comment, number, at)
876 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
877 )
878 .bind(&[
879 new_id("act", now_ms).into(),
880 repo_id.into(),
881 item.alert_id.into(),
882 item.action.into(),
883 optional(item.actor),
884 optional(item.reason.map(DismissReason::as_str)),
885 optional(item.comment),
886 item.number.map_or(JsValue::NULL, JsValue::from),
887 at.as_str().into(),
888 ])
889 })
890 .collect::<Result<Vec<_>>>()?;
891 self.db.batch(statements).await?;
892 Ok(())
893 }
894
895 /// What happened to a repository's alerts, newest first.
896 pub async fn activity(&self, repo_id: &str, limit: u32) -> Result<Vec<AlertActivity>> {
897 Ok(self
898 .db
899 .prepare("SELECT * FROM alert_activity WHERE repo_id = ? ORDER BY at DESC, id DESC LIMIT ?")
900 .bind(&[repo_id.into(), limit.into()])?
901 .all()
902 .await?
903 .results::<ActivityRow>()?
904 .into_iter()
905 .map(AlertActivity::from)
906 .collect())
907 }
908
909 // --- Security updates ---------------------------------------------------
910
911 pub async fn updates(&self, repo_id: &str) -> Result<Vec<UpdateRow>> {
912 self.db
913 .prepare("SELECT * FROM updates WHERE repo_id = ?")
914 .bind(&[repo_id.into()])?
915 .all()
916 .await?
917 .results::<UpdateRow>()
918 }
919
920 pub async fn update(&self, repo_id: &str, ecosystem: &str, package: &str) -> Result<Option<UpdateRow>> {
921 self.db
922 .prepare("SELECT * FROM updates WHERE repo_id = ? AND ecosystem = ? AND package = ?")
923 .bind(&[repo_id.into(), ecosystem.into(), package.into()])?
924 .first::<UpdateRow>(None)
925 .await
926 }
927
928 pub async fn update_by_branch(&self, repo_id: &str, branch: &str) -> Result<Option<UpdateRow>> {
929 self.db
930 .prepare("SELECT * FROM updates WHERE repo_id = ? AND branch = ?")
931 .bind(&[repo_id.into(), branch.into()])?
932 .first::<UpdateRow>(None)
933 .await
934 }
935
936 pub async fn update_by_pull(&self, repo_id: &str, pull: u32) -> Result<Option<UpdateRow>> {
937 self.db
938 .prepare("SELECT * FROM updates WHERE repo_id = ? AND pull = ?")
939 .bind(&[repo_id.into(), pull.into()])?
940 .first::<UpdateRow>(None)
941 .await
942 }
943
944 /// Updates asked of a sandbox before `before` that never pushed.
945 pub async fn stalled_updates(&self, before: &str, limit: u32) -> Result<Vec<UpdateRow>> {
946 self.db
947 .prepare("SELECT * FROM updates WHERE state = 'requested' AND updated_at < ? ORDER BY updated_at LIMIT ?")
948 .bind(&[before.into(), limit.into()])?
949 .all()
950 .await?
951 .results::<UpdateRow>()
952 }
953
954 /// Asks for a security update: a sandbox is making the change on
955 /// `branch`. An open pull request for an older version is kept until
956 /// the new one opens, which supersedes it.
957 pub async fn request_update(&self, repo_id: &str, ecosystem: &str, package: &str, target: &str, branch: &str) -> Result<()> {
958 let now = now();
959 self.db
960 .prepare(
961 "INSERT INTO updates (repo_id, ecosystem, package, target, state, branch, requested_at, updated_at)
962 VALUES (?1, ?2, ?3, ?4, 'requested', ?5, ?6, ?6)
963 ON CONFLICT (repo_id, ecosystem, package) DO UPDATE SET
964 target = ?4, state = 'requested', branch = ?5, error = NULL, issue = NULL,
965 pull = CASE WHEN updates.state = 'open' THEN updates.pull ELSE NULL END,
966 requested_at = ?6, updated_at = ?6",
967 )
968 .bind(&[repo_id.into(), ecosystem.into(), package.into(), target.into(), branch.into(), now.into()])?
969 .run()
970 .await?;
971 Ok(())
972 }
973
974 pub async fn set_update(
975 &self,
976 row: &UpdateRow,
977 state: UpdateState,
978 pull: Option<u32>,
979 issue: Option<u32>,
980 error: Option<&str>,
981 ) -> Result<()> {
982 self.db
983 .prepare(
984 "UPDATE updates SET state = ?, pull = ?, issue = ?, error = ?, updated_at = ?
985 WHERE repo_id = ? AND ecosystem = ? AND package = ?",
986 )
987 .bind(&[
988 state.as_str().into(),
989 pull.map_or(JsValue::NULL, JsValue::from),
990 issue.map_or(JsValue::NULL, JsValue::from),
991 optional(error),
992 now().into(),
993 row.repo_id.as_str().into(),
994 row.ecosystem.as_str().into(),
995 row.package.as_str().into(),
996 ])?
997 .run()
998 .await?;
999 Ok(())
1000 }
1001
1002 // --- Pushes scanned after they landed -----------------------------------
1003
1004 pub async fn add_push_scan(&self, repo_id: &str, git_ref: &str, head: &str, base: Option<&str>, pusher: Option<&str>) -> Result<String> {
1005 let id = new_id("psc", now_ms());
1006 self.db
1007 .prepare(
1008 "INSERT INTO push_scans (id, repo_id, git_ref, head, base, pusher, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)",
1009 )
1010 .bind(&[id.as_str().into(), repo_id.into(), git_ref.into(), head.into(), optional(base), optional(pusher), now().into()])?
1011 .run()
1012 .await?;
1013 Ok(id)
1014 }
1015
1016 pub async fn push_scan(&self, id: &str) -> Result<Option<PushScanRow>> {
1017 self.db
1018 .prepare("SELECT * FROM push_scans WHERE id = ?")
1019 .bind(&[id.into()])?
1020 .first::<PushScanRow>(None)
1021 .await
1022 }
1023
1024 pub async fn pending_push_scans(&self, limit: u32) -> Result<Vec<PushScanRow>> {
1025 self.db
1026 .prepare("SELECT * FROM push_scans WHERE state = 'pending' ORDER BY created_at LIMIT ?")
1027 .bind(&[limit.into()])?
1028 .all()
1029 .await?
1030 .results::<PushScanRow>()
1031 }
1032
1033 /// Where a push's scan stands after a page: its next page, or done.
1034 pub async fn advance_push_scan(&self, id: &str, next: Option<&str>, commits: u32, found: u32) -> Result<()> {
1035 let done = next.is_none();
1036 self.db
1037 .prepare(
1038 "UPDATE push_scans SET cursor = ?, state = ?, commits = commits + ?, pages = pages + 1, found = found + ?,
1039 finished_at = CASE WHEN ? THEN ? ELSE finished_at END
1040 WHERE id = ?",
1041 )
1042 .bind(&[
1043 optional(next),
1044 if done { "done" } else { "pending" }.into(),
1045 commits.into(),
1046 found.into(),
1047 done.into(),
1048 now().into(),
1049 id.into(),
1050 ])?
1051 .run()
1052 .await?;
1053 Ok(())
1054 }
1055
1056 // --- Upgrades -----------------------------------------------------------
1057
1058 /// The issue opened for a package before security updates, if any.
1059 pub async fn upgrade(&self, repo_id: &str, ecosystem: &str, package: &str) -> Result<Option<UpgradeRow>> {
1060 self.db
1061 .prepare("SELECT number FROM upgrades WHERE repo_id = ? AND ecosystem = ? AND package = ?")
1062 .bind(&[repo_id.into(), ecosystem.into(), package.into()])?
1063 .first::<UpgradeRow>(None)
1064 .await
1065 }
1066
1067 // --- Advisories and usage -----------------------------------------------
1068
1069 /// OSV's record of a vulnerability, if one was fetched after `after`.
1070 pub async fn advisory(&self, osv_id: &str, after: &str) -> Result<Option<serde_json::Value>> {
1071 let row = self
1072 .db
1073 .prepare("SELECT body, fetched_at FROM advisories WHERE osv_id = ?")
1074 .bind(&[osv_id.into()])?
1075 .first::<BodyRow>(None)
1076 .await?;
1077 Ok(row
1078 .filter(|row| row.fetched_at.as_str() >= after)
1079 .and_then(|row| serde_json::from_str(&row.body).ok()))
1080 }
1081
1082 pub async fn keep_advisory(&self, osv_id: &str, body: &serde_json::Value) -> Result<()> {
1083 self.db
1084 .prepare(
1085 "INSERT INTO advisories (osv_id, body, fetched_at) VALUES (?1, ?2, ?3)
1086 ON CONFLICT (osv_id) DO UPDATE SET body = ?2, fetched_at = ?3",
1087 )
1088 .bind(&[osv_id.into(), body.to_string().into(), now().into()])?
1089 .run()
1090 .await?;
1091 Ok(())
1092 }
1093
1094 /// Adds to a workspace's scanning this month; returns what the month
1095 /// has cost so far, in millionths of a dollar.
1096 pub async fn meter(&self, workspace: &str, reads: u32, commits: u32, osv_queries: u32, cost_micros: i64) -> Result<i64> {
1097 let month = now()[..7].to_owned();
1098 self.db
1099 .prepare(
1100 "INSERT INTO usage (workspace, month, reads, commits, osv_queries, cost_micros) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
1101 ON CONFLICT (workspace, month) DO UPDATE SET
1102 reads = reads + ?3, commits = commits + ?4, osv_queries = osv_queries + ?5,
1103 cost_micros = cost_micros + ?6",
1104 )
1105 .bind(&[
1106 workspace.into(),
1107 month.as_str().into(),
1108 reads.into(),
1109 commits.into(),
1110 osv_queries.into(),
1111 JsValue::from(cost_micros as f64),
1112 ])?
1113 .run()
1114 .await?;
1115 Ok(self
1116 .db
1117 .prepare("SELECT cost_micros AS n FROM usage WHERE workspace = ? AND month = ?")
1118 .bind(&[workspace.into(), month.into()])?
1119 .first::<NumberRow>(None)
1120 .await?
1121 .map_or(0, |row| row.n))
1122 }
1123}
1124
1125/// Secret alerts counted by where they stand, from (status, looks like a
1126/// test value, how many).
1127pub fn secret_counts<'a>(rows: impl Iterator<Item = (&'a str, bool, u32)>) -> SecretCounts {
1128 let mut counts = SecretCounts::default();
1129 for (status, test, n) in rows {
1130 match (SecretStatus::parse(status), test) {
1131 (Some(SecretStatus::Open | SecretStatus::Blocked), true) => counts.test_values += n,
1132 (Some(SecretStatus::Open), false) => counts.open += n,
1133 (Some(SecretStatus::Blocked), false) => counts.blocked += n,
1134 (Some(SecretStatus::Allowed), _) => counts.dismissed += n,
1135 (Some(SecretStatus::Resolved), _) => counts.fixed += n,
1136 (None, _) => {}
1137 }
1138 }
1139 counts
1140}
1141
1142#[cfg(test)]
1143mod tests {
1144 use super::*;
1145
1146 #[test]
1147 fn only_real_secrets_in_history_are_open() {
1148 let counts = secret_counts(
1149 [("open", false, 2), ("open", true, 3), ("blocked", false, 1), ("blocked", true, 4), ("allowed", true, 5), ("resolved", false, 6)]
1150 .into_iter(),
1151 );
1152 assert_eq!((counts.open, counts.blocked, counts.test_values, counts.dismissed, counts.fixed), (2, 1, 7, 5, 6));
1153 }
1154}