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

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