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/deps.rs

461 lines20,976 bytesCodeBlame
1//! Dependencies: reading a repository's lockfiles, asking OSV about every
2//! package in them, and opening one upgrade issue per vulnerable package
3//! for a g1t agent to land through the normal pull request flow.
4
5use std::collections::{BTreeMap, BTreeSet, HashMap};
6
7use g1t_contracts::repos::RepoPath;
8use g1t_contracts::security::{FindLockfilesArgs, Lockfiles, VulnStatus, Vulnerability};
9use g1t_contracts::time::rfc3339;
10use g1t_contracts::work::{AddCommentArgs, Issue, IssueDetail, IssueReason, OpenIssueArgs, QueueIssueArgs, State, ViewArgs};
11use g1t_contracts::{Outcome, User};
12use g1t_kit::now_ms;
13use g1t_scan::lockfiles::{Lockfile, Package, still_locked_check, test_command};
14use g1t_scan::osv::{self, Advisory, Severity};
15use serde_json::{Value, json};
16use worker::{Fetch, Headers, Method, Request, RequestInit, Result};
17
18use crate::Security;
19use crate::store::{RepoRow, VulnRow};
20
21/// OSV's records are fetched again after this long.
22const ADVISORY_MAX_AGE_MS: u64 = 7 * 24 * 60 * 60 * 1000;
23/// Records fetched per scan, at most; the rest wait for the next one.
24const MAX_ADVISORY_FETCHES: usize = 150;
25/// Upgrade issues opened per scan, most severe first.
26const MAX_NEW_ISSUES: usize = 8;
27/// CPU one call to OSV takes, sending it and reading its answer, in
28/// milliseconds (an estimate, rounded up). OSV itself is free, and a
29/// Worker's outgoing requests are not charged.
30const CPU_MS_PER_OSV_CALL: f64 = 2.0;
31
32/// What a dependency check cost g1t, in millionths of a dollar, rounded
33/// up, at the prices in `history`: the CPU of its OSV calls, and the rows
34/// it writes (an advisory kept for each call at most, each vulnerability
35/// found, where the check stands and the month's usage).
36pub fn dependency_check_cost(calls: u32, found: usize) -> i64 {
37 use crate::history::{MICROS_PER_CPU_MS, MICROS_PER_ROW_WRITTEN};
38 let cpu = f64::from(calls) * CPU_MS_PER_OSV_CALL * MICROS_PER_CPU_MS;
39 let rows = (calls as usize + found + 2) as f64 * MICROS_PER_ROW_WRITTEN;
40 (cpu + rows).ceil() as i64
41}
42
43async fn osv_call(method: Method, url: &str, body: Option<&Value>) -> Result<Option<Value>> {
44 let headers = Headers::new();
45 headers.set("user-agent", "g1t (+https://g1t.sh)")?;
46 headers.set("accept", "application/json")?;
47 if body.is_some() {
48 headers.set("content-type", "application/json")?;
49 }
50 let mut init = RequestInit::new();
51 init.with_method(method).with_headers(headers);
52 if let Some(body) = body {
53 init.with_body(Some(body.to_string().into()));
54 }
55 let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?;
56 if response.status_code() == 404 {
57 return Ok(None);
58 }
59 if !(200..300).contains(&response.status_code()) {
60 return Err(worker::Error::RustError(format!("OSV answered {}", response.status_code())));
61 }
62 Ok(Some(response.json().await?))
63}
64
65/// One package in one lockfile.
66struct Located {
67 package: Package,
68 lockfile: Lockfile,
69 path: String,
70}
71
72/// Every package the lockfiles resolve. A directory with a `go.mod` is read
73/// from it rather than from its `go.sum`, which lists versions not built.
74fn packages(files: &Lockfiles) -> Vec<Located> {
75 let go_mods: BTreeSet<&str> = files
76 .files
77 .iter()
78 .filter(|file| file.path.ends_with("go.mod"))
79 .map(|file| file.path.trim_end_matches("go.mod"))
80 .collect();
81 let mut located = Vec::new();
82 for file in &files.files {
83 let Some(lockfile) = Lockfile::for_path(&file.path) else { continue };
84 if lockfile == Lockfile::GoSum && go_mods.contains(file.path.trim_end_matches("go.sum")) {
85 continue;
86 }
87 for package in lockfile.parse(&file.text) {
88 located.push(Located { package, lockfile, path: file.path.clone() });
89 }
90 }
91 located
92}
93
94fn directory(path: &str) -> &str {
95 path.rsplit_once('/').map_or("", |(directory, _)| directory)
96}
97
98impl Security {
99 /// The ids of the vulnerabilities affecting each package, from OSV.
100 async fn query_osv(&self, packages: &[Package]) -> Result<(Vec<Vec<String>>, u32)> {
101 let mut ids = Vec::with_capacity(packages.len());
102 let mut calls = 0;
103 for (body, chunk) in osv::batch_bodies(packages).iter().zip(packages.chunks(osv::MAX_BATCH)) {
104 calls += 1;
105 let answer = osv_call(Method::Post, osv::QUERY_BATCH_URL, Some(body)).await?.unwrap_or(Value::Null);
106 let (mut found, more) = osv::read_batch(&answer, chunk.len());
107 // A package with many advisories is paged; fetch the rest.
108 for (index, token) in more.into_iter().take(20) {
109 let package = &chunk[index];
110 let query = json!({
111 "package": {"name": package.name, "ecosystem": package.ecosystem.osv()},
112 "version": package.version,
113 "page_token": token,
114 });
115 calls += 1;
116 if let Some(page) = osv_call(Method::Post, "https://api.osv.dev/v1/query", Some(&query)).await? {
117 found[index].extend(
118 page["vulns"].as_array().into_iter().flatten().filter_map(|v| v["id"].as_str().map(str::to_owned)),
119 );
120 }
121 }
122 ids.extend(found);
123 }
124 Ok((ids, calls))
125 }
126
127 /// OSV's record of each id, from the cache when it is fresh.
128 async fn advisories(&self, ids: &BTreeSet<String>) -> Result<(HashMap<String, Value>, u32)> {
129 let fresh_after = rfc3339(now_ms().saturating_sub(ADVISORY_MAX_AGE_MS));
130 let mut records = HashMap::new();
131 let mut fetched = 0u32;
132 for id in ids {
133 if let Some(record) = self.store.advisory(id, &fresh_after).await? {
134 records.insert(id.clone(), record);
135 continue;
136 }
137 if fetched as usize >= MAX_ADVISORY_FETCHES {
138 continue;
139 }
140 fetched += 1;
141 if let Some(record) = osv_call(Method::Get, &osv::vuln_url(id), None).await? {
142 self.store.keep_advisory(id, &record).await?;
143 records.insert(id.clone(), record);
144 }
145 }
146 Ok((records, fetched))
147 }
148
149 /// Reads a repository's dependencies, records which are vulnerable,
150 /// and opens upgrade issues for those with a fix. Returns what went
151 /// wrong, for the Security page, if anything did.
152 pub async fn scan_dependencies(&self, repo: &RepoRow) -> Result<Option<String>> {
153 let files: Lockfiles = g1t_kit::call(&self.repos, "find_lockfiles", &FindLockfilesArgs { repo_id: repo.repo_id.clone() }).await?;
154 let paths: Vec<String> = files.files.iter().map(|file| file.path.clone()).collect();
155 let located = packages(&files);
156 let unique: Vec<Package> = located.iter().map(|l| l.package.clone()).collect::<BTreeSet<_>>().into_iter().collect();
157 let outcome = async {
158 let (ids, calls) = self.query_osv(&unique).await?;
159 let by_package: HashMap<&Package, &Vec<String>> = unique.iter().zip(ids.iter()).collect();
160 let wanted: BTreeSet<String> = ids.iter().flatten().cloned().collect();
161 let (records, fetched) = self.advisories(&wanted).await?;
162 let mut found = Vec::new();
163 for item in &located {
164 for id in by_package.get(&item.package).into_iter().flat_map(|ids| ids.iter()) {
165 let Some(advisory) = records.get(id).and_then(|record| osv::read_vuln(record, &item.package)) else {
166 continue;
167 };
168 found.push(vulnerability(&repo.repo_id, item, &advisory));
169 }
170 }
171 Ok::<_, worker::Error>((found, calls + fetched))
172 }
173 .await;
174 let (found, calls) = match outcome {
175 Ok(result) => result,
176 Err(error) => {
177 let problem = format!("The dependencies could not be checked: {error}");
178 self.store
179 .set_dependencies_scanned(&repo.repo_id, files.commit.as_deref(), &paths, Some(&problem))
180 .await?;
181 return Ok(Some(problem));
182 }
183 };
184 self.store.replace_vulnerabilities(&repo.repo_id, &found).await?;
185 self.store.set_dependencies_scanned(&repo.repo_id, files.commit.as_deref(), &paths, None).await?;
186 self.meter(&repo.namespace, 0, 0, calls, dependency_check_cost(calls, found.len())).await?;
187 // No upgrade issues on an archived (read-only) or deleted repository.
188 if repo.upkeep != 0 && self.active(&repo.repo_id).await? {
189 self.open_upgrades(repo, &located).await?;
190 }
191 Ok(None)
192 }
193
194 /// One issue per vulnerable package that has a fix and no open issue
195 /// for it, most severe first, with a g1t agent put on the first and the
196 /// rest queued for one.
197 async fn open_upgrades(&self, repo: &RepoRow, located: &[Located]) -> Result<()> {
198 let open = self.store.open_vulnerabilities(&repo.repo_id).await?;
199 let mut by_package: BTreeMap<(String, String), Vec<&VulnRow>> = BTreeMap::new();
200 for vuln in &open {
201 if vuln.fixed_version.is_some() {
202 by_package.entry((vuln.ecosystem.clone(), vuln.package.clone())).or_default().push(vuln);
203 }
204 }
205 let mut groups: Vec<((String, String), Vec<&VulnRow>)> = by_package.into_iter().collect();
206 groups.sort_by_key(|(_, vulns)| std::cmp::Reverse(vulns.iter().map(|v| Severity::parse(&v.severity)).max()));
207 let Some(actor) = self.workspace_actor(&repo.namespace).await? else {
208 return Ok(());
209 };
210 let path = RepoPath { namespace: repo.namespace.clone(), name: repo.name.clone() };
211 let mut opened = 0;
212 let mut agents: Option<std::result::Result<(), String>> = None;
213 for ((ecosystem, package), vulns) in groups {
214 if opened >= MAX_NEW_ISSUES {
215 break;
216 }
217 if let Some(existing) = self.store.upgrade(&repo.repo_id, &ecosystem, &package).await? {
218 match self.issue(&actor, &path, existing.number as u32).await? {
219 // Already being fixed.
220 Some(issue) if issue.state == State::Open => continue,
221 // Someone decided not to; respect it until they reopen it.
222 Some(issue) if issue.reason == Some(IssueReason::NotPlanned) => continue,
223 _ => {}
224 }
225 }
226 let target = osv::upgrade_target(vulns.iter().filter_map(|v| v.fixed_version.as_deref()))
227 .unwrap_or_default();
228 let mut advisories: Vec<&str> = vulns.iter().map(|v| v.advisory.as_str()).collect();
229 advisories.sort();
230 advisories.dedup();
231 let named = match advisories.as_slice() {
232 [one] => (*one).to_owned(),
233 [first, second] => format!("{first}, {second}"),
234 [first, rest @ ..] => format!("{first} and {} more", rest.len()),
235 [] => "a known vulnerability".to_owned(),
236 };
237 let title: String = format!("Upgrade {package} to {target}: fixes {named}").chars().take(200).collect();
238 let (body, checks) = issue_text(&ecosystem, &package, &target, &vulns, located);
239 let issue: Outcome<Issue> = g1t_kit::call(
240 &self.work,
241 "open_issue",
242 &OpenIssueArgs {
243 actor: actor.clone(),
244 repo: path.clone(),
245 title,
246 body,
247 labels: vec!["dependencies".to_owned(), "security".to_owned()],
248 checks,
249 },
250 )
251 .await?;
252 let Outcome::Ok(issue) = issue else { continue };
253 opened += 1;
254 // The first upgrade starts an agent at once, which also says
255 // whether this workspace can run agents; the rest wait in the
256 // queue, which starts them as the repository has room.
257 let (assigned, note) = match &agents {
258 None => {
259 let started: Outcome<Value> = g1t_kit::call(
260 &self.runner,
261 "run",
262 &json!({ "actor": actor, "repo": path, "issue": issue.number }),
263 )
264 .await?;
265 let result = match started {
266 Outcome::Ok(_) => Ok(()),
267 Outcome::Fail(refused) => Err(refused.message),
268 };
269 agents = Some(result.clone());
270 match result {
271 Ok(()) => (true, None),
272 Err(reason) => (false, Some(reason)),
273 }
274 }
275 Some(Ok(())) => {
276 let queued: Outcome<bool> = g1t_kit::call(
277 &self.work,
278 "queue_issue",
279 &QueueIssueArgs { actor: actor.clone(), repo: path.clone(), number: issue.number, queued: true },
280 )
281 .await?;
282 (matches!(queued, Outcome::Ok(true)), None)
283 }
284 Some(Err(reason)) => (false, Some(reason.clone())),
285 };
286 if let Some(reason) = &note {
287 self.comment(&actor, &path, issue.number, format!(
288 "g1t could not put an agent on this upgrade: {reason}\n\nAssign it to g1t-agent once agents can run here, or upgrade it by hand."
289 ))
290 .await?;
291 }
292 self.store
293 .record_upgrade(&repo.repo_id, &ecosystem, &package, issue.number, &target, assigned, note.as_deref())
294 .await?;
295 }
296 Ok(())
297 }
298
299 async fn issue(&self, actor: &User, repo: &RepoPath, number: u32) -> Result<Option<Issue>> {
300 let found: Outcome<IssueDetail> = g1t_kit::call(
301 &self.work,
302 "get_issue",
303 &ViewArgs { repo: repo.clone(), number, viewer: Some(actor.clone()), after_seq: 0 },
304 )
305 .await?;
306 Ok(found.into_result().ok().map(|detail| detail.issue))
307 }
308
309 async fn comment(&self, actor: &User, repo: &RepoPath, number: u32, body: String) -> Result<()> {
310 let _: Outcome<Value> = g1t_kit::call(
311 &self.work,
312 "add_comment",
313 &AddCommentArgs {
314 actor: actor.clone(),
315 repo: repo.clone(),
316 number,
317 body,
318 path: None,
319 line: None,
320 verdict: None,
321 },
322 )
323 .await?;
324 Ok(())
325 }
326}
327
328fn vulnerability(repo_id: &str, item: &Located, advisory: &Advisory) -> Vulnerability {
329 Vulnerability {
330 id: String::new(),
331 repo_id: repo_id.to_owned(),
332 ecosystem: item.package.ecosystem.osv().to_owned(),
333 package: item.package.name.clone(),
334 version: item.package.version.clone(),
335 manifest: item.path.clone(),
336 advisory: advisory.display_id.clone(),
337 osv_id: advisory.id.clone(),
338 summary: advisory.summary.clone(),
339 severity: advisory.severity.as_str().to_owned(),
340 fixed_version: advisory.fixed.clone(),
341 status: VulnStatus::Open,
342 issue: None,
343 found_at: String::new(),
344 fixed_at: None,
345 }
346}
347
348/// The issue's body, written for the agent that takes it as much as for a
349/// person, and its acceptance checks: the project's tests, and that no
350/// lockfile still resolves a vulnerable version.
351fn issue_text(ecosystem: &str, package: &str, target: &str, vulns: &[&VulnRow], located: &[Located]) -> (String, Vec<String>) {
352 let mut body = format!(
353 "`{package}` ({ecosystem}) has known vulnerabilities with a fix in **{target}**. Upgrade it to {target} or later \
354 everywhere it is locked, keeping other changes to what the upgrade needs.\n\n\
355 | Advisory | Severity | Affected | Fixed in | Summary |\n| --- | --- | --- | --- | --- |\n"
356 );
357 let mut seen = BTreeSet::new();
358 for vuln in vulns {
359 if !seen.insert((vuln.advisory.clone(), vuln.version.clone())) {
360 continue;
361 }
362 body.push_str(&format!(
363 "| [{}]({}) | {} | {} | {} | {} |\n",
364 vuln.advisory,
365 osv::page_url(&vuln.osv_id),
366 vuln.severity,
367 vuln.version,
368 vuln.fixed_version.as_deref().unwrap_or("none yet"),
369 vuln.summary.replace('|', "\\|").replace('\n', " "),
370 ));
371 }
372 let mut checks = Vec::new();
373 let mut manifests = BTreeSet::new();
374 let mut tests = BTreeSet::new();
375 for vuln in vulns {
376 let Some(item) = located
377 .iter()
378 .find(|item| item.path == vuln.manifest && item.package.name == vuln.package && item.package.version == vuln.version)
379 else {
380 continue;
381 };
382 if manifests.insert((item.path.clone(), item.package.version.clone())) {
383 checks.push(still_locked_check(item.lockfile, &item.path, &item.package.name, &item.package.version));
384 }
385 if let Some(test) = test_command(item.lockfile, directory(&item.path)) {
386 tests.insert(test);
387 }
388 }
389 let locked: Vec<String> = manifests.iter().map(|(path, version)| format!("`{path}` ({version})")).collect();
390 body.push_str(&format!("\nLocked in: {}.\n", locked.join(", ")));
391 body.push_str(
392 "\nThe acceptance checks pass once no lockfile resolves a vulnerable version and the tests still pass. \
393 If the fix needs a major upgrade that breaks the build, change the code that depends on it in the same pull request.\n\n\
394 ---\n_Opened by g1t's dependency upkeep. Turn it off for this project on its Security page._",
395 );
396 checks.extend(tests);
397 (body, checks)
398}
399
400#[cfg(test)]
401mod tests {
402 use super::*;
403 use g1t_contracts::security::LockfileText;
404
405 #[test]
406 fn scans_cost_their_cpu_and_the_rows_they_write() {
407 // 10 OSV calls and 3 vulnerabilities: 0.4 of CPU, 15 rows.
408 assert_eq!(dependency_check_cost(10, 3), 16);
409 assert_eq!(dependency_check_cost(0, 0), 2);
410 // A page of 25 commits that read 100 objects and found nothing:
411 // 10 of CPU and 2 rows. The old placeholder charged 100.
412 assert_eq!(crate::history::history_page_cost(100, 0), 12);
413 assert_eq!(crate::history::history_page_cost(0, 1), 3);
414 }
415
416 #[test]
417 fn go_sum_is_skipped_beside_go_mod() {
418 let files = Lockfiles {
419 commit: None,
420 files: vec![
421 LockfileText { path: "go.mod".into(), text: "require golang.org/x/net v0.7.0\n".into() },
422 LockfileText { path: "go.sum".into(), text: "golang.org/x/net v0.1.0 h1:x=\n".into() },
423 LockfileText { path: "tools/go.sum".into(), text: "golang.org/x/text v0.3.0 h1:x=\n".into() },
424 ],
425 };
426 let found: Vec<String> = packages(&files).iter().map(|l| format!("{}:{}", l.path, l.package.version)).collect();
427 assert_eq!(found, ["go.mod:v0.7.0", "tools/go.sum:v0.3.0"]);
428 }
429
430 #[test]
431 fn the_issue_names_the_advisories_and_checks_the_lockfile() {
432 let row = VulnRow {
433 id: "vul_1".into(),
434 repo_id: "rep_1".into(),
435 ecosystem: "npm".into(),
436 package: "lodash".into(),
437 version: "4.17.20".into(),
438 manifest: "web/package-lock.json".into(),
439 osv_id: "GHSA-35jh-r3h4-6jhm".into(),
440 advisory: "GHSA-35jh-r3h4-6jhm".into(),
441 summary: "Command Injection in lodash".into(),
442 severity: "high".into(),
443 fixed_version: Some("4.17.21".into()),
444 status: "open".into(),
445 found_at: "2026-10-04T00:00:00Z".into(),
446 fixed_at: None,
447 number: None,
448 };
449 let located = vec![Located {
450 package: Package { ecosystem: g1t_scan::lockfiles::Ecosystem::Npm, name: "lodash".into(), version: "4.17.20".into() },
451 lockfile: Lockfile::PackageLock,
452 path: "web/package-lock.json".into(),
453 }];
454 let (body, checks) = issue_text("npm", "lodash", "4.17.21", &[&row], &located);
455 assert!(body.contains("[GHSA-35jh-r3h4-6jhm](https://osv.dev/vulnerability/GHSA-35jh-r3h4-6jhm) | high | 4.17.20 | 4.17.21"));
456 assert!(body.contains("`web/package-lock.json` (4.17.20)"));
457 assert_eq!(checks.len(), 2);
458 assert!(checks[0].contains("node_modules/lodash") && checks[0].contains("'web/package-lock.json'"));
459 assert_eq!(checks[1], "cd 'web' && npm ci && npm test --if-present");
460 }
461}