flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/security/src/store.rs

620 lines22,651 bytesCodeBlame
1//! The security service's tables in D1.
2
3use g1t_contracts::security::{
4 NewSecret, ScanState, SecretFinding, SecretStatus, SeverityCounts, VulnStatus, Vulnerability,
5};
6use g1t_contracts::time::rfc3339;
7use g1t_contracts::new_id;
8use g1t_kit::now_ms;
9use g1t_scan::secrets::SecretKind;
10use serde::Deserialize;
11use worker::wasm_bindgen::JsValue;
12use worker::{D1Database, Result};
13
14pub fn optional(value: Option<&str>) -> JsValue {
15 value.map_or(JsValue::NULL, JsValue::from)
16}
17
18pub fn now() -> String {
19 rfc3339(now_ms())
20}
21
22#[derive(Clone, Deserialize)]
23pub struct RepoRow {
24 pub repo_id: String,
25 pub namespace: String,
26 pub name: String,
27 pub upkeep: i64,
28 pub history: String,
29 pub history_cursor: Option<String>,
30 pub history_commits: i64,
31 pub history_finished_at: Option<String>,
32 pub deps_scanned_at: Option<String>,
33 pub deps_error: Option<String>,
34 pub lockfiles: String,
35}
36
37impl RepoRow {
38 pub fn scan_state(&self) -> ScanState {
39 ScanState {
40 history: self.history.clone(),
41 commits_scanned: self.history_commits.max(0) as u32,
42 history_finished_at: self.history_finished_at.clone(),
43 dependencies_scanned_at: self.deps_scanned_at.clone(),
44 dependencies_error: self.deps_error.clone(),
45 lockfiles: serde_json::from_str(&self.lockfiles).unwrap_or_default(),
46 }
47 }
48}
49
50#[derive(Deserialize)]
51struct SecretRow {
52 id: String,
53 repo_id: String,
54 kind: String,
55 path: String,
56 line: i64,
57 commit_hash: String,
58 preview: String,
59 status: String,
60 source: String,
61 found_by: Option<String>,
62 found_at: String,
63 decided_by: Option<String>,
64 reason: Option<String>,
65 decided_at: Option<String>,
66}
67
68impl From<SecretRow> for SecretFinding {
69 fn from(row: SecretRow) -> Self {
70 SecretFinding {
71 label: SecretKind::parse(&row.kind).map_or("a secret", |kind| kind.label()).to_owned(),
72 id: row.id,
73 repo_id: row.repo_id,
74 kind: row.kind,
75 path: row.path,
76 line: row.line.max(0) as u32,
77 commit: row.commit_hash,
78 preview: row.preview,
79 status: SecretStatus::parse(&row.status).unwrap_or(SecretStatus::Open),
80 source: row.source,
81 found_by: row.found_by,
82 found_at: row.found_at,
83 decided_by: row.decided_by,
84 reason: row.reason,
85 decided_at: row.decided_at,
86 }
87 }
88}
89
90#[derive(Deserialize)]
91pub struct VulnRow {
92 pub id: String,
93 pub repo_id: String,
94 pub ecosystem: String,
95 pub package: String,
96 pub version: String,
97 pub manifest: String,
98 pub osv_id: String,
99 pub advisory: String,
100 pub summary: String,
101 pub severity: String,
102 pub fixed_version: Option<String>,
103 pub status: String,
104 pub found_at: String,
105 pub fixed_at: Option<String>,
106 pub number: Option<i64>,
107}
108
109impl From<VulnRow> for Vulnerability {
110 fn from(row: VulnRow) -> Self {
111 Vulnerability {
112 id: row.id,
113 repo_id: row.repo_id,
114 ecosystem: row.ecosystem,
115 package: row.package,
116 version: row.version,
117 manifest: row.manifest,
118 advisory: row.advisory,
119 osv_id: row.osv_id,
120 summary: row.summary,
121 severity: row.severity,
122 fixed_version: row.fixed_version,
123 status: if row.status == "fixed" { VulnStatus::Fixed } else { VulnStatus::Open },
124 issue: row.number.map(|number| number as u32),
125 found_at: row.found_at,
126 fixed_at: row.fixed_at,
127 }
128 }
129}
130
131#[derive(Deserialize)]
132pub struct UpgradeRow {
133 pub number: i64,
134}
135
136#[derive(Deserialize)]
137struct CountRow {
138 severity: String,
139 n: i64,
140}
141
142#[derive(Deserialize)]
143struct NumberRow {
144 n: i64,
145}
146
147#[derive(Deserialize)]
148struct FingerprintRow {
149 fingerprint: String,
150 id: String,
151 status: String,
152}
153
154#[derive(Deserialize)]
155struct BodyRow {
156 body: String,
157 fetched_at: String,
158}
159
160pub struct Store {
161 pub db: D1Database,
162}
163
164const VULN_COLUMNS: &str = "v.id, v.repo_id, v.ecosystem, v.package, v.version, v.manifest, v.osv_id, v.advisory,
165 v.summary, v.severity, v.fixed_version, v.status, v.found_at, v.fixed_at, u.number AS number";
166
167impl Store {
168 pub async fn repo(&self, repo_id: &str) -> Result<Option<RepoRow>> {
169 self.db
170 .prepare("SELECT * FROM repos WHERE repo_id = ?")
171 .bind(&[repo_id.into()])?
172 .first::<RepoRow>(None)
173 .await
174 }
175
176 /// Records a repository the first time it is seen, and keeps its
177 /// address current. Returns its row.
178 pub async fn register(&self, repo_id: &str, namespace: &str, name: &str) -> Result<RepoRow> {
179 self.db
180 .prepare(
181 "INSERT INTO repos (repo_id, namespace, name, created_at) VALUES (?1, ?2, ?3, ?4)
182 ON CONFLICT (repo_id) DO UPDATE SET namespace = ?2, name = ?3",
183 )
184 .bind(&[repo_id.into(), namespace.into(), name.into(), now().into()])?
185 .run()
186 .await?;
187 self.repo(repo_id)
188 .await?
189 .ok_or_else(|| worker::Error::RustError("the repository was not recorded".into()))
190 }
191
192 pub async fn rename_namespace(&self, stale: &[String], current: &str) -> Result<()> {
193 for slug in stale {
194 self.db
195 .prepare("UPDATE repos SET namespace = ? WHERE namespace = ?")
196 .bind(&[current.into(), slug.as_str().into()])?
197 .run()
198 .await?;
199 }
200 Ok(())
201 }
202
203 pub async fn in_namespace(&self, namespace: &str) -> Result<Vec<RepoRow>> {
204 self.db
205 .prepare("SELECT * FROM repos WHERE namespace = ? ORDER BY name")
206 .bind(&[namespace.into()])?
207 .all()
208 .await?
209 .results::<RepoRow>()
210 }
211
212 /// Repositories whose history still has to be scanned, oldest first.
213 pub async fn unfinished_histories(&self, limit: u32) -> Result<Vec<RepoRow>> {
214 self.db
215 .prepare("SELECT * FROM repos WHERE history IN ('pending', 'running', 'stopped') ORDER BY created_at LIMIT ?")
216 .bind(&[limit.into()])?
217 .all()
218 .await?
219 .results::<RepoRow>()
220 }
221
222 /// Repositories whose dependencies were last read before `before`.
223 pub async fn stale_dependencies(&self, before: &str, limit: u32) -> Result<Vec<RepoRow>> {
224 self.db
225 .prepare(
226 "SELECT * FROM repos WHERE deps_scanned_at IS NULL OR deps_scanned_at < ?
227 ORDER BY deps_scanned_at LIMIT ?",
228 )
229 .bind(&[before.into(), limit.into()])?
230 .all()
231 .await?
232 .results::<RepoRow>()
233 }
234
235 pub async fn set_history(&self, repo_id: &str, state: &str, cursor: Option<&str>, commits: u32) -> Result<()> {
236 let finished = (state == "done").then(now);
237 self.db
238 .prepare(
239 "UPDATE repos SET history = ?, history_cursor = ?, history_commits = history_commits + ?,
240 history_finished_at = COALESCE(?, history_finished_at)
241 WHERE repo_id = ?",
242 )
243 .bind(&[state.into(), optional(cursor), commits.into(), optional(finished.as_deref()), repo_id.into()])?
244 .run()
245 .await?;
246 Ok(())
247 }
248
249 pub async fn restart_history(&self, repo_id: &str) -> Result<()> {
250 self.db
251 .prepare(
252 "UPDATE repos SET history = 'pending', history_cursor = NULL, history_commits = 0,
253 history_finished_at = NULL WHERE repo_id = ?",
254 )
255 .bind(&[repo_id.into()])?
256 .run()
257 .await?;
258 Ok(())
259 }
260
261 pub async fn set_dependencies_scanned(
262 &self,
263 repo_id: &str,
264 commit: Option<&str>,
265 lockfiles: &[String],
266 error: Option<&str>,
267 ) -> Result<()> {
268 self.db
269 .prepare("UPDATE repos SET deps_scanned_at = ?, deps_commit = ?, lockfiles = ?, deps_error = ? WHERE repo_id = ?")
270 .bind(&[now().into(), optional(commit), serde_json::to_string(lockfiles)?.into(), optional(error), repo_id.into()])?
271 .run()
272 .await?;
273 Ok(())
274 }
275
276 pub async fn set_upkeep(&self, repo_id: &str, enabled: bool, by: &str) -> Result<()> {
277 self.db
278 .prepare("UPDATE repos SET upkeep = ?, upkeep_by = ?, upkeep_at = ? WHERE repo_id = ?")
279 .bind(&[i32::from(enabled).into(), by.into(), now().into(), repo_id.into()])?
280 .run()
281 .await?;
282 Ok(())
283 }
284
285 // --- Secrets ------------------------------------------------------------
286
287 pub async fn secrets(&self, repo_id: &str) -> Result<Vec<SecretFinding>> {
288 let rows = self
289 .db
290 .prepare(
291 "SELECT * FROM secrets WHERE repo_id = ?
292 ORDER BY CASE status WHEN 'open' THEN 0 WHEN 'blocked' THEN 1 WHEN 'allowed' THEN 2 ELSE 3 END,
293 found_at DESC LIMIT 500",
294 )
295 .bind(&[repo_id.into()])?
296 .all()
297 .await?
298 .results::<SecretRow>()?;
299 Ok(rows.into_iter().map(SecretFinding::from).collect())
300 }
301
302 pub async fn secret(&self, repo_id: &str, id: &str) -> Result<Option<SecretFinding>> {
303 Ok(self
304 .db
305 .prepare("SELECT * FROM secrets WHERE repo_id = ? AND id = ?")
306 .bind(&[repo_id.into(), id.into()])?
307 .first::<SecretRow>(None)
308 .await?
309 .map(SecretFinding::from))
310 }
311
312 /// The findings already known for these fingerprints: fingerprint, id
313 /// and status.
314 pub async fn known(&self, repo_id: &str, fingerprints: &[String]) -> Result<Vec<(String, String, String)>> {
315 let mut known = Vec::new();
316 for chunk in fingerprints.chunks(90) {
317 let marks = vec!["?"; chunk.len()].join(", ");
318 let mut binds: Vec<JsValue> = vec![repo_id.into()];
319 binds.extend(chunk.iter().map(|fingerprint| JsValue::from(fingerprint.as_str())));
320 let rows = self
321 .db
322 .prepare(format!("SELECT fingerprint, id, status FROM secrets WHERE repo_id = ? AND fingerprint IN ({marks})"))
323 .bind(&binds)?
324 .all()
325 .await?
326 .results::<FingerprintRow>()?;
327 known.extend(rows.into_iter().map(|row| (row.fingerprint, row.id, row.status)));
328 }
329 Ok(known)
330 }
331
332 /// Records secrets not seen before in this repository, with `status`.
333 /// A secret already known keeps its record and its decision.
334 pub async fn add_secrets(
335 &self,
336 repo_id: &str,
337 secrets: &[NewSecret],
338 status: SecretStatus,
339 source: &str,
340 found_by: Option<&str>,
341 ) -> Result<()> {
342 if secrets.is_empty() {
343 return Ok(());
344 }
345 let now_ms = now_ms();
346 let found_at = rfc3339(now_ms);
347 let statements = secrets
348 .iter()
349 .map(|secret| {
350 self.db
351 .prepare(
352 "INSERT INTO secrets (id, repo_id, fingerprint, kind, path, line, commit_hash, preview,
353 status, source, found_by, found_at)
354 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
355 ON CONFLICT (repo_id, fingerprint) DO NOTHING",
356 )
357 .bind(&[
358 new_id("sec", now_ms).into(),
359 repo_id.into(),
360 secret.fingerprint.as_str().into(),
361 secret.kind.as_str().into(),
362 secret.path.as_str().into(),
363 secret.line.into(),
364 secret.commit.as_str().into(),
365 secret.preview.as_str().into(),
366 status.as_str().into(),
367 source.into(),
368 optional(found_by),
369 found_at.as_str().into(),
370 ])
371 })
372 .collect::<Result<Vec<_>>>()?;
373 self.db.batch(statements).await?;
374 Ok(())
375 }
376
377 /// A secret found in history that was only ever blocked has landed
378 /// after all (it was allowed, then pushed): it is open now unless
379 /// someone allowed it.
380 pub async fn landed(&self, repo_id: &str, fingerprints: &[String]) -> Result<()> {
381 for chunk in fingerprints.chunks(90) {
382 let marks = vec!["?"; chunk.len()].join(", ");
383 let mut binds: Vec<JsValue> = vec![repo_id.into()];
384 binds.extend(chunk.iter().map(|fingerprint| JsValue::from(fingerprint.as_str())));
385 self.db
386 .prepare(format!(
387 "UPDATE secrets SET status = 'open', source = 'history'
388 WHERE repo_id = ? AND status = 'blocked' AND fingerprint IN ({marks})"
389 ))
390 .bind(&binds)?
391 .run()
392 .await?;
393 }
394 Ok(())
395 }
396
397 pub async fn decide(&self, repo_id: &str, id: &str, status: SecretStatus, by: &str, reason: Option<&str>) -> Result<()> {
398 let reopened = status == SecretStatus::Open;
399 self.db
400 .prepare("UPDATE secrets SET status = ?, decided_by = ?, reason = ?, decided_at = ? WHERE repo_id = ? AND id = ?")
401 .bind(&[
402 status.as_str().into(),
403 optional((!reopened).then_some(by)),
404 optional(reason.filter(|_| !reopened)),
405 optional((!reopened).then(now).as_deref()),
406 repo_id.into(),
407 id.into(),
408 ])?
409 .run()
410 .await?;
411 Ok(())
412 }
413
414 // --- Vulnerabilities ----------------------------------------------------
415
416 pub async fn vulnerabilities(&self, repo_id: &str) -> Result<Vec<Vulnerability>> {
417 let rows = self
418 .db
419 .prepare(format!(
420 "SELECT {VULN_COLUMNS} FROM vulnerabilities v
421 LEFT JOIN upgrades u ON u.repo_id = v.repo_id AND u.ecosystem = v.ecosystem AND u.package = v.package
422 WHERE v.repo_id = ?
423 ORDER BY CASE v.status WHEN 'open' THEN 0 ELSE 1 END,
424 CASE v.severity WHEN 'critical' THEN 0 WHEN 'high' THEN 1 WHEN 'medium' THEN 2 WHEN 'low' THEN 3 ELSE 4 END,
425 v.package LIMIT 1000"
426 ))
427 .bind(&[repo_id.into()])?
428 .all()
429 .await?
430 .results::<VulnRow>()?;
431 Ok(rows.into_iter().map(Vulnerability::from).collect())
432 }
433
434 pub async fn open_vulnerabilities(&self, repo_id: &str) -> Result<Vec<VulnRow>> {
435 self.db
436 .prepare(format!(
437 "SELECT {VULN_COLUMNS} FROM vulnerabilities v
438 LEFT JOIN upgrades u ON u.repo_id = v.repo_id AND u.ecosystem = v.ecosystem AND u.package = v.package
439 WHERE v.repo_id = ? AND v.status = 'open'"
440 ))
441 .bind(&[repo_id.into()])?
442 .all()
443 .await?
444 .results::<VulnRow>()
445 }
446
447 /// Replaces what is known about a repository's dependencies with what a
448 /// scan found: new findings open, findings no longer true fixed, and
449 /// findings that came back open again.
450 pub async fn replace_vulnerabilities(&self, repo_id: &str, found: &[Vulnerability]) -> Result<()> {
451 let now_ms = now_ms();
452 let now = rfc3339(now_ms);
453 let mut statements = vec![self
454 .db
455 .prepare("UPDATE vulnerabilities SET status = 'fixed', fixed_at = ? WHERE repo_id = ? AND status = 'open'")
456 .bind(&[now.as_str().into(), repo_id.into()])?];
457 for vuln in found {
458 statements.push(
459 self.db
460 .prepare(
461 "INSERT INTO vulnerabilities (id, repo_id, ecosystem, package, version, manifest, osv_id, advisory,
462 summary, severity, fixed_version, status, found_at)
463 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'open', ?)
464 ON CONFLICT (repo_id, ecosystem, package, version, manifest, osv_id) DO UPDATE SET
465 status = 'open', fixed_at = NULL, advisory = excluded.advisory, summary = excluded.summary,
466 severity = excluded.severity, fixed_version = excluded.fixed_version",
467 )
468 .bind(&[
469 new_id("vul", now_ms).into(),
470 repo_id.into(),
471 vuln.ecosystem.as_str().into(),
472 vuln.package.as_str().into(),
473 vuln.version.as_str().into(),
474 vuln.manifest.as_str().into(),
475 vuln.osv_id.as_str().into(),
476 vuln.advisory.as_str().into(),
477 vuln.summary.as_str().into(),
478 vuln.severity.as_str().into(),
479 optional(vuln.fixed_version.as_deref()),
480 now.as_str().into(),
481 ])?,
482 );
483 }
484 // D1 runs a batch as one transaction, so readers never see the
485 // moment between marking everything fixed and opening it again.
486 self.db.batch(statements).await?;
487 Ok(())
488 }
489
490 pub async fn counts(&self, repo_id: &str) -> Result<(SeverityCounts, u32, u32)> {
491 let rows = self
492 .db
493 .prepare("SELECT severity, count(*) AS n FROM vulnerabilities WHERE repo_id = ? AND status = 'open' GROUP BY severity")
494 .bind(&[repo_id.into()])?
495 .all()
496 .await?
497 .results::<CountRow>()?;
498 let mut counts = SeverityCounts::default();
499 let mut vulns = 0;
500 for row in rows {
501 let n = row.n.max(0) as u32;
502 vulns += n;
503 match row.severity.as_str() {
504 "critical" => counts.critical += n,
505 "high" => counts.high += n,
506 "medium" => counts.medium += n,
507 "low" => counts.low += n,
508 _ => counts.unknown += n,
509 }
510 }
511 let secrets = self
512 .db
513 .prepare("SELECT count(*) AS n FROM secrets WHERE repo_id = ? AND status IN ('open', 'blocked')")
514 .bind(&[repo_id.into()])?
515 .first::<NumberRow>(None)
516 .await?
517 .map_or(0, |row| row.n.max(0) as u32);
518 counts.critical += secrets;
519 Ok((counts, secrets, vulns))
520 }
521
522 // --- Upgrades -----------------------------------------------------------
523
524 pub async fn upgrade(&self, repo_id: &str, ecosystem: &str, package: &str) -> Result<Option<UpgradeRow>> {
525 self.db
526 .prepare("SELECT number FROM upgrades WHERE repo_id = ? AND ecosystem = ? AND package = ?")
527 .bind(&[repo_id.into(), ecosystem.into(), package.into()])?
528 .first::<UpgradeRow>(None)
529 .await
530 }
531
532 pub async fn record_upgrade(
533 &self,
534 repo_id: &str,
535 ecosystem: &str,
536 package: &str,
537 number: u32,
538 target: &str,
539 assigned: bool,
540 note: Option<&str>,
541 ) -> Result<()> {
542 self.db
543 .prepare(
544 "INSERT INTO upgrades (repo_id, ecosystem, package, number, target, opened_at, assigned, note)
545 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
546 ON CONFLICT (repo_id, ecosystem, package) DO UPDATE SET
547 number = ?4, target = ?5, opened_at = ?6, assigned = ?7, note = ?8",
548 )
549 .bind(&[
550 repo_id.into(),
551 ecosystem.into(),
552 package.into(),
553 number.into(),
554 target.into(),
555 now().into(),
556 i32::from(assigned).into(),
557 optional(note),
558 ])?
559 .run()
560 .await?;
561 Ok(())
562 }
563
564 // --- Advisories and usage -----------------------------------------------
565
566 /// OSV's record of a vulnerability, if one was fetched after `after`.
567 pub async fn advisory(&self, osv_id: &str, after: &str) -> Result<Option<serde_json::Value>> {
568 let row = self
569 .db
570 .prepare("SELECT body, fetched_at FROM advisories WHERE osv_id = ?")
571 .bind(&[osv_id.into()])?
572 .first::<BodyRow>(None)
573 .await?;
574 Ok(row
575 .filter(|row| row.fetched_at.as_str() >= after)
576 .and_then(|row| serde_json::from_str(&row.body).ok()))
577 }
578
579 pub async fn keep_advisory(&self, osv_id: &str, body: &serde_json::Value) -> Result<()> {
580 self.db
581 .prepare(
582 "INSERT INTO advisories (osv_id, body, fetched_at) VALUES (?1, ?2, ?3)
583 ON CONFLICT (osv_id) DO UPDATE SET body = ?2, fetched_at = ?3",
584 )
585 .bind(&[osv_id.into(), body.to_string().into(), now().into()])?
586 .run()
587 .await?;
588 Ok(())
589 }
590
591 /// Adds to a workspace's scanning this month; returns what the month
592 /// has cost so far, in millionths of a dollar.
593 pub async fn meter(&self, workspace: &str, reads: u32, commits: u32, osv_queries: u32, cost_micros: i64) -> Result<i64> {
594 let month = now()[..7].to_owned();
595 self.db
596 .prepare(
597 "INSERT INTO usage (workspace, month, reads, commits, osv_queries, cost_micros) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
598 ON CONFLICT (workspace, month) DO UPDATE SET
599 reads = reads + ?3, commits = commits + ?4, osv_queries = osv_queries + ?5,
600 cost_micros = cost_micros + ?6",
601 )
602 .bind(&[
603 workspace.into(),
604 month.as_str().into(),
605 reads.into(),
606 commits.into(),
607 osv_queries.into(),
608 JsValue::from(cost_micros as f64),
609 ])?
610 .run()
611 .await?;
612 Ok(self
613 .db
614 .prepare("SELECT cost_micros AS n FROM usage WHERE workspace = ? AND month = ?")
615 .bind(&[workspace.into(), month.into()])?
616 .first::<NumberRow>(None)
617 .await?
618 .map_or(0, |row| row.n))
619 }
620}