Skip to content
1,164 linesCodeBlameRaw
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 /// Every vulnerability alert of a repository: open, fixed and dismissed.
773 pub async fn all_vulnerabilities(&self, repo_id: &str) -> Result<Vec<VulnRow>> {
774 self.db
775 .prepare(format!("SELECT {VULN_COLUMNS} FROM vulnerabilities v {VULN_JOINS} WHERE v.repo_id = ?"))
776 .bind(&[repo_id.into()])?
777 .all()
778 .await?
779 .results::<VulnRow>()
780 }
781
782 /// Replaces what is known about a repository's dependencies with what a
783 /// scan found: new findings open, findings no longer true fixed, and
784 /// findings that came back open again, or dismissed again when someone
785 /// had dismissed them.
786 pub async fn replace_vulnerabilities(&self, repo_id: &str, found: &[Vulnerability]) -> Result<()> {
787 let now_ms = now_ms();
788 let now = rfc3339(now_ms);
789 let mut statements = vec![self
790 .db
791 .prepare(
792 "UPDATE vulnerabilities SET status = 'fixed', fixed_at = ? WHERE repo_id = ? AND status IN ('open', 'dismissed')",
793 )
794 .bind(&[now.as_str().into(), repo_id.into()])?];
795 for vuln in found {
796 statements.push(
797 self.db
798 .prepare(
799 "INSERT INTO vulnerabilities (id, repo_id, ecosystem, package, version, manifest, osv_id, advisory,
800 summary, severity, fixed_version, status, found_at)
801 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'open', ?)
802 ON CONFLICT (repo_id, ecosystem, package, version, manifest, osv_id) DO UPDATE SET
803 status = CASE WHEN vulnerabilities.dismiss_reason IS NULL THEN 'open' ELSE 'dismissed' END,
804 fixed_at = NULL, advisory = excluded.advisory, summary = excluded.summary,
805 severity = excluded.severity, fixed_version = excluded.fixed_version",
806 )
807 .bind(&[
808 new_id("vul", now_ms).into(),
809 repo_id.into(),
810 vuln.ecosystem.as_str().into(),
811 vuln.package.as_str().into(),
812 vuln.version.as_str().into(),
813 vuln.manifest.as_str().into(),
814 vuln.osv_id.as_str().into(),
815 vuln.advisory.as_str().into(),
816 vuln.summary.as_str().into(),
817 vuln.severity.as_str().into(),
818 optional(vuln.fixed_version.as_deref()),
819 now.as_str().into(),
820 ])?,
821 );
822 }
823 // D1 runs a batch as one transaction, so readers never see the
824 // moment between marking everything fixed and opening it again.
825 self.db.batch(statements).await?;
826 Ok(())
827 }
828
829 pub async fn counts(&self, repo_id: &str) -> Result<(SeverityCounts, u32, u32)> {
830 let rows = self
831 .db
832 .prepare("SELECT severity, count(*) AS n FROM vulnerabilities WHERE repo_id = ? AND status = 'open' GROUP BY severity")
833 .bind(&[repo_id.into()])?
834 .all()
835 .await?
836 .results::<CountRow>()?;
837 let mut counts = SeverityCounts::default();
838 let mut vulns = 0;
839 for row in rows {
840 let n = row.n.max(0) as u32;
841 vulns += n;
842 match row.severity.as_str() {
843 "critical" => counts.critical += n,
844 "high" => counts.high += n,
845 "medium" => counts.medium += n,
846 "low" => counts.low += n,
847 _ => counts.unknown += n,
848 }
849 }
850 let secrets = self.secret_counts(repo_id).await?;
851 // A secret in the history that looks real is the worst thing a
852 // repository can hold. One stopped at a push never landed, and a
853 // likely test value is no danger: neither is counted.
854 counts.critical += secrets.open;
855 Ok((counts, secrets.open + secrets.blocked, vulns))
856 }
857
858 pub async fn secret_counts(&self, repo_id: &str) -> Result<SecretCounts> {
859 let rows = self
860 .db
861 .prepare(
862 "SELECT status, (test_value IS NOT NULL) AS test, count(*) AS n FROM secrets WHERE repo_id = ?
863 GROUP BY status, test",
864 )
865 .bind(&[repo_id.into()])?
866 .all()
867 .await?
868 .results::<SecretCountRow>()?;
869 Ok(secret_counts(rows.iter().map(|row| (row.status.as_str(), row.test != 0, row.n.max(0) as u32))))
870 }
871
872 // --- Activity -----------------------------------------------------------
873
874 pub async fn record(&self, repo_id: &str, activity: &[Activity<'_>]) -> Result<()> {
875 if activity.is_empty() {
876 return Ok(());
877 }
878 let now_ms = now_ms();
879 let at = rfc3339(now_ms);
880 let statements = activity
881 .iter()
882 .map(|item| {
883 self.db
884 .prepare(
885 "INSERT INTO alert_activity (id, repo_id, alert_id, action, actor, reason, comment, number, at)
886 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
887 )
888 .bind(&[
889 new_id("act", now_ms).into(),
890 repo_id.into(),
891 item.alert_id.into(),
892 item.action.into(),
893 optional(item.actor),
894 optional(item.reason.map(DismissReason::as_str)),
895 optional(item.comment),
896 item.number.map_or(JsValue::NULL, JsValue::from),
897 at.as_str().into(),
898 ])
899 })
900 .collect::<Result<Vec<_>>>()?;
901 self.db.batch(statements).await?;
902 Ok(())
903 }
904
905 /// What happened to a repository's alerts, newest first.
906 pub async fn activity(&self, repo_id: &str, limit: u32) -> Result<Vec<AlertActivity>> {
907 Ok(self
908 .db
909 .prepare("SELECT * FROM alert_activity WHERE repo_id = ? ORDER BY at DESC, id DESC LIMIT ?")
910 .bind(&[repo_id.into(), limit.into()])?
911 .all()
912 .await?
913 .results::<ActivityRow>()?
914 .into_iter()
915 .map(AlertActivity::from)
916 .collect())
917 }
918
919 // --- Security updates ---------------------------------------------------
920
921 pub async fn updates(&self, repo_id: &str) -> Result<Vec<UpdateRow>> {
922 self.db
923 .prepare("SELECT * FROM updates WHERE repo_id = ?")
924 .bind(&[repo_id.into()])?
925 .all()
926 .await?
927 .results::<UpdateRow>()
928 }
929
930 pub async fn update(&self, repo_id: &str, ecosystem: &str, package: &str) -> Result<Option<UpdateRow>> {
931 self.db
932 .prepare("SELECT * FROM updates WHERE repo_id = ? AND ecosystem = ? AND package = ?")
933 .bind(&[repo_id.into(), ecosystem.into(), package.into()])?
934 .first::<UpdateRow>(None)
935 .await
936 }
937
938 pub async fn update_by_branch(&self, repo_id: &str, branch: &str) -> Result<Option<UpdateRow>> {
939 self.db
940 .prepare("SELECT * FROM updates WHERE repo_id = ? AND branch = ?")
941 .bind(&[repo_id.into(), branch.into()])?
942 .first::<UpdateRow>(None)
943 .await
944 }
945
946 pub async fn update_by_pull(&self, repo_id: &str, pull: u32) -> Result<Option<UpdateRow>> {
947 self.db
948 .prepare("SELECT * FROM updates WHERE repo_id = ? AND pull = ?")
949 .bind(&[repo_id.into(), pull.into()])?
950 .first::<UpdateRow>(None)
951 .await
952 }
953
954 /// Updates asked of a sandbox before `before` that never pushed.
955 pub async fn stalled_updates(&self, before: &str, limit: u32) -> Result<Vec<UpdateRow>> {
956 self.db
957 .prepare("SELECT * FROM updates WHERE state = 'requested' AND updated_at < ? ORDER BY updated_at LIMIT ?")
958 .bind(&[before.into(), limit.into()])?
959 .all()
960 .await?
961 .results::<UpdateRow>()
962 }
963
964 /// Asks for a security update: a sandbox is making the change on
965 /// `branch`. An open pull request for an older version is kept until
966 /// the new one opens, which supersedes it.
967 pub async fn request_update(&self, repo_id: &str, ecosystem: &str, package: &str, target: &str, branch: &str) -> Result<()> {
968 let now = now();
969 self.db
970 .prepare(
971 "INSERT INTO updates (repo_id, ecosystem, package, target, state, branch, requested_at, updated_at)
972 VALUES (?1, ?2, ?3, ?4, 'requested', ?5, ?6, ?6)
973 ON CONFLICT (repo_id, ecosystem, package) DO UPDATE SET
974 target = ?4, state = 'requested', branch = ?5, error = NULL, issue = NULL,
975 pull = CASE WHEN updates.state = 'open' THEN updates.pull ELSE NULL END,
976 requested_at = ?6, updated_at = ?6",
977 )
978 .bind(&[repo_id.into(), ecosystem.into(), package.into(), target.into(), branch.into(), now.into()])?
979 .run()
980 .await?;
981 Ok(())
982 }
983
984 pub async fn set_update(
985 &self,
986 row: &UpdateRow,
987 state: UpdateState,
988 pull: Option<u32>,
989 issue: Option<u32>,
990 error: Option<&str>,
991 ) -> Result<()> {
992 self.db
993 .prepare(
994 "UPDATE updates SET state = ?, pull = ?, issue = ?, error = ?, updated_at = ?
995 WHERE repo_id = ? AND ecosystem = ? AND package = ?",
996 )
997 .bind(&[
998 state.as_str().into(),
999 pull.map_or(JsValue::NULL, JsValue::from),
1000 issue.map_or(JsValue::NULL, JsValue::from),
1001 optional(error),
1002 now().into(),
1003 row.repo_id.as_str().into(),
1004 row.ecosystem.as_str().into(),
1005 row.package.as_str().into(),
1006 ])?
1007 .run()
1008 .await?;
1009 Ok(())
1010 }
1011
1012 // --- Pushes scanned after they landed -----------------------------------
1013
1014 pub async fn add_push_scan(&self, repo_id: &str, git_ref: &str, head: &str, base: Option<&str>, pusher: Option<&str>) -> Result<String> {
1015 let id = new_id("psc", now_ms());
1016 self.db
1017 .prepare(
1018 "INSERT INTO push_scans (id, repo_id, git_ref, head, base, pusher, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)",
1019 )
1020 .bind(&[id.as_str().into(), repo_id.into(), git_ref.into(), head.into(), optional(base), optional(pusher), now().into()])?
1021 .run()
1022 .await?;
1023 Ok(id)
1024 }
1025
1026 pub async fn push_scan(&self, id: &str) -> Result<Option<PushScanRow>> {
1027 self.db
1028 .prepare("SELECT * FROM push_scans WHERE id = ?")
1029 .bind(&[id.into()])?
1030 .first::<PushScanRow>(None)
1031 .await
1032 }
1033
1034 pub async fn pending_push_scans(&self, limit: u32) -> Result<Vec<PushScanRow>> {
1035 self.db
1036 .prepare("SELECT * FROM push_scans WHERE state = 'pending' ORDER BY created_at LIMIT ?")
1037 .bind(&[limit.into()])?
1038 .all()
1039 .await?
1040 .results::<PushScanRow>()
1041 }
1042
1043 /// Where a push's scan stands after a page: its next page, or done.
1044 pub async fn advance_push_scan(&self, id: &str, next: Option<&str>, commits: u32, found: u32) -> Result<()> {
1045 let done = next.is_none();
1046 self.db
1047 .prepare(
1048 "UPDATE push_scans SET cursor = ?, state = ?, commits = commits + ?, pages = pages + 1, found = found + ?,
1049 finished_at = CASE WHEN ? THEN ? ELSE finished_at END
1050 WHERE id = ?",
1051 )
1052 .bind(&[
1053 optional(next),
1054 if done { "done" } else { "pending" }.into(),
1055 commits.into(),
1056 found.into(),
1057 done.into(),
1058 now().into(),
1059 id.into(),
1060 ])?
1061 .run()
1062 .await?;
1063 Ok(())
1064 }
1065
1066 // --- Upgrades -----------------------------------------------------------
1067
1068 /// The issue opened for a package before security updates, if any.
1069 pub async fn upgrade(&self, repo_id: &str, ecosystem: &str, package: &str) -> Result<Option<UpgradeRow>> {
1070 self.db
1071 .prepare("SELECT number FROM upgrades WHERE repo_id = ? AND ecosystem = ? AND package = ?")
1072 .bind(&[repo_id.into(), ecosystem.into(), package.into()])?
1073 .first::<UpgradeRow>(None)
1074 .await
1075 }
1076
1077 // --- Advisories and usage -----------------------------------------------
1078
1079 /// OSV's record of a vulnerability, if one was fetched after `after`.
1080 pub async fn advisory(&self, osv_id: &str, after: &str) -> Result<Option<serde_json::Value>> {
1081 let row = self
1082 .db
1083 .prepare("SELECT body, fetched_at FROM advisories WHERE osv_id = ?")
1084 .bind(&[osv_id.into()])?
1085 .first::<BodyRow>(None)
1086 .await?;
1087 Ok(row
1088 .filter(|row| row.fetched_at.as_str() >= after)
1089 .and_then(|row| serde_json::from_str(&row.body).ok()))
1090 }
1091
1092 pub async fn keep_advisory(&self, osv_id: &str, body: &serde_json::Value) -> Result<()> {
1093 self.db
1094 .prepare(
1095 "INSERT INTO advisories (osv_id, body, fetched_at) VALUES (?1, ?2, ?3)
1096 ON CONFLICT (osv_id) DO UPDATE SET body = ?2, fetched_at = ?3",
1097 )
1098 .bind(&[osv_id.into(), body.to_string().into(), now().into()])?
1099 .run()
1100 .await?;
1101 Ok(())
1102 }
1103
1104 /// Adds to a workspace's scanning this month; returns what the month
1105 /// has cost so far, in millionths of a dollar.
1106 pub async fn meter(&self, workspace: &str, reads: u32, commits: u32, osv_queries: u32, cost_micros: i64) -> Result<i64> {
1107 let month = now()[..7].to_owned();
1108 self.db
1109 .prepare(
1110 "INSERT INTO usage (workspace, month, reads, commits, osv_queries, cost_micros) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
1111 ON CONFLICT (workspace, month) DO UPDATE SET
1112 reads = reads + ?3, commits = commits + ?4, osv_queries = osv_queries + ?5,
1113 cost_micros = cost_micros + ?6",
1114 )
1115 .bind(&[
1116 workspace.into(),
1117 month.as_str().into(),
1118 reads.into(),
1119 commits.into(),
1120 osv_queries.into(),
1121 JsValue::from(cost_micros as f64),
1122 ])?
1123 .run()
1124 .await?;
1125 Ok(self
1126 .db
1127 .prepare("SELECT cost_micros AS n FROM usage WHERE workspace = ? AND month = ?")
1128 .bind(&[workspace.into(), month.into()])?
1129 .first::<NumberRow>(None)
1130 .await?
1131 .map_or(0, |row| row.n))
1132 }
1133}
1134
1135/// Secret alerts counted by where they stand, from (status, looks like a
1136/// test value, how many).
1137pub fn secret_counts<'a>(rows: impl Iterator<Item = (&'a str, bool, u32)>) -> SecretCounts {
1138 let mut counts = SecretCounts::default();
1139 for (status, test, n) in rows {
1140 match (SecretStatus::parse(status), test) {
1141 (Some(SecretStatus::Open | SecretStatus::Blocked), true) => counts.test_values += n,
1142 (Some(SecretStatus::Open), false) => counts.open += n,
1143 (Some(SecretStatus::Blocked), false) => counts.blocked += n,
1144 (Some(SecretStatus::Allowed), _) => counts.dismissed += n,
1145 (Some(SecretStatus::Resolved), _) => counts.fixed += n,
1146 (None, _) => {}
1147 }
1148 }
1149 counts
1150}
1151
1152#[cfg(test)]
1153mod tests {
1154 use super::*;
1155
1156 #[test]
1157 fn only_real_secrets_in_history_are_open() {
1158 let counts = secret_counts(
1159 [("open", false, 2), ("open", true, 3), ("blocked", false, 1), ("blocked", true, 4), ("allowed", true, 5), ("resolved", false, 6)]
1160 .into_iter(),
1161 );
1162 assert_eq!((counts.open, counts.blocked, counts.test_values, counts.dismissed, counts.fixed), (2, 1, 7, 5, 6));
1163 }
1164}