flagon-io/g1t

public

Git for AI scale: a forge for thousands of agents working on the same code at once.

g1t/services/security/src/deps.rs

460 lines20,860 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API1//! 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;
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put27/// 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;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API31
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put32/// 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
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API43async 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?;
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put186 self.meter(&repo.namespace, 0, 0, calls, dependency_check_cost(calls, found.len())).await?;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API187 if repo.upkeep != 0 {
188 self.open_upgrades(repo, &located).await?;
189 }
190 Ok(None)
191 }
192
193 /// One issue per vulnerable package that has a fix and no open issue
194 /// for it, most severe first, with a g1t agent put on the first and the
195 /// rest queued for one.
196 async fn open_upgrades(&self, repo: &RepoRow, located: &[Located]) -> Result<()> {
197 let open = self.store.open_vulnerabilities(&repo.repo_id).await?;
198 let mut by_package: BTreeMap<(String, String), Vec<&VulnRow>> = BTreeMap::new();
199 for vuln in &open {
200 if vuln.fixed_version.is_some() {
201 by_package.entry((vuln.ecosystem.clone(), vuln.package.clone())).or_default().push(vuln);
202 }
203 }
204 let mut groups: Vec<((String, String), Vec<&VulnRow>)> = by_package.into_iter().collect();
205 groups.sort_by_key(|(_, vulns)| std::cmp::Reverse(vulns.iter().map(|v| Severity::parse(&v.severity)).max()));
206 let Some(actor) = self.workspace_actor(&repo.namespace).await? else {
207 return Ok(());
208 };
209 let path = RepoPath { namespace: repo.namespace.clone(), name: repo.name.clone() };
210 let mut opened = 0;
211 let mut agents: Option<std::result::Result<(), String>> = None;
212 for ((ecosystem, package), vulns) in groups {
213 if opened >= MAX_NEW_ISSUES {
214 break;
215 }
216 if let Some(existing) = self.store.upgrade(&repo.repo_id, &ecosystem, &package).await? {
217 match self.issue(&actor, &path, existing.number as u32).await? {
218 // Already being fixed.
219 Some(issue) if issue.state == State::Open => continue,
220 // Someone decided not to; respect it until they reopen it.
221 Some(issue) if issue.reason == Some(IssueReason::NotPlanned) => continue,
222 _ => {}
223 }
224 }
225 let target = osv::upgrade_target(vulns.iter().filter_map(|v| v.fixed_version.as_deref()))
226 .unwrap_or_default();
227 let mut advisories: Vec<&str> = vulns.iter().map(|v| v.advisory.as_str()).collect();
228 advisories.sort();
229 advisories.dedup();
230 let named = match advisories.as_slice() {
231 [one] => (*one).to_owned(),
232 [first, second] => format!("{first}, {second}"),
233 [first, rest @ ..] => format!("{first} and {} more", rest.len()),
234 [] => "a known vulnerability".to_owned(),
235 };
236 let title: String = format!("Upgrade {package} to {target}: fixes {named}").chars().take(200).collect();
237 let (body, checks) = issue_text(&ecosystem, &package, &target, &vulns, located);
238 let issue: Outcome<Issue> = g1t_kit::call(
239 &self.work,
240 "open_issue",
241 &OpenIssueArgs {
242 actor: actor.clone(),
243 repo: path.clone(),
244 title,
245 body,
246 labels: vec!["dependencies".to_owned(), "security".to_owned()],
247 checks,
248 },
249 )
250 .await?;
251 let Outcome::Ok(issue) = issue else { continue };
252 opened += 1;
253 // The first upgrade starts an agent at once, which also says
254 // whether this workspace can run agents; the rest wait in the
255 // queue, which starts them as the repository has room.
256 let (assigned, note) = match &agents {
257 None => {
258 let started: Outcome<Value> = g1t_kit::call(
259 &self.runner,
260 "run",
261 &json!({ "actor": actor, "repo": path, "issue": issue.number }),
262 )
263 .await?;
264 let result = match started {
265 Outcome::Ok(_) => Ok(()),
266 Outcome::Fail(refused) => Err(refused.message),
267 };
268 agents = Some(result.clone());
269 match result {
270 Ok(()) => (true, None),
271 Err(reason) => (false, Some(reason)),
272 }
273 }
274 Some(Ok(())) => {
275 let queued: Outcome<bool> = g1t_kit::call(
276 &self.work,
277 "queue_issue",
278 &QueueIssueArgs { actor: actor.clone(), repo: path.clone(), number: issue.number, queued: true },
279 )
280 .await?;
281 (matches!(queued, Outcome::Ok(true)), None)
282 }
283 Some(Err(reason)) => (false, Some(reason.clone())),
284 };
285 if let Some(reason) = &note {
286 self.comment(&actor, &path, issue.number, format!(
287 "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."
288 ))
289 .await?;
290 }
291 self.store
292 .record_upgrade(&repo.repo_id, &ecosystem, &package, issue.number, &target, assigned, note.as_deref())
293 .await?;
294 }
295 Ok(())
296 }
297
298 async fn issue(&self, actor: &User, repo: &RepoPath, number: u32) -> Result<Option<Issue>> {
299 let found: Outcome<IssueDetail> = g1t_kit::call(
300 &self.work,
301 "get_issue",
302 &ViewArgs { repo: repo.clone(), number, viewer: Some(actor.clone()), after_seq: 0 },
303 )
304 .await?;
305 Ok(found.into_result().ok().map(|detail| detail.issue))
306 }
307
308 async fn comment(&self, actor: &User, repo: &RepoPath, number: u32, body: String) -> Result<()> {
309 let _: Outcome<Value> = g1t_kit::call(
310 &self.work,
311 "add_comment",
312 &AddCommentArgs {
313 actor: actor.clone(),
314 repo: repo.clone(),
315 number,
316 body,
317 path: None,
318 line: None,
319 verdict: None,
320 },
321 )
322 .await?;
323 Ok(())
324 }
325}
326
327fn vulnerability(repo_id: &str, item: &Located, advisory: &Advisory) -> Vulnerability {
328 Vulnerability {
329 id: String::new(),
330 repo_id: repo_id.to_owned(),
331 ecosystem: item.package.ecosystem.osv().to_owned(),
332 package: item.package.name.clone(),
333 version: item.package.version.clone(),
334 manifest: item.path.clone(),
335 advisory: advisory.display_id.clone(),
336 osv_id: advisory.id.clone(),
337 summary: advisory.summary.clone(),
338 severity: advisory.severity.as_str().to_owned(),
339 fixed_version: advisory.fixed.clone(),
340 status: VulnStatus::Open,
341 issue: None,
342 found_at: String::new(),
343 fixed_at: None,
344 }
345}
346
347/// The issue's body, written for the agent that takes it as much as for a
348/// person, and its acceptance checks: the project's tests, and that no
349/// lockfile still resolves a vulnerable version.
350fn issue_text(ecosystem: &str, package: &str, target: &str, vulns: &[&VulnRow], located: &[Located]) -> (String, Vec<String>) {
351 let mut body = format!(
352 "`{package}` ({ecosystem}) has known vulnerabilities with a fix in **{target}**. Upgrade it to {target} or later \
353 everywhere it is locked, keeping other changes to what the upgrade needs.\n\n\
354 | Advisory | Severity | Affected | Fixed in | Summary |\n| --- | --- | --- | --- | --- |\n"
355 );
356 let mut seen = BTreeSet::new();
357 for vuln in vulns {
358 if !seen.insert((vuln.advisory.clone(), vuln.version.clone())) {
359 continue;
360 }
361 body.push_str(&format!(
362 "| [{}]({}) | {} | {} | {} | {} |\n",
363 vuln.advisory,
364 osv::page_url(&vuln.osv_id),
365 vuln.severity,
366 vuln.version,
367 vuln.fixed_version.as_deref().unwrap_or("none yet"),
368 vuln.summary.replace('|', "\\|").replace('\n', " "),
369 ));
370 }
371 let mut checks = Vec::new();
372 let mut manifests = BTreeSet::new();
373 let mut tests = BTreeSet::new();
374 for vuln in vulns {
375 let Some(item) = located
376 .iter()
377 .find(|item| item.path == vuln.manifest && item.package.name == vuln.package && item.package.version == vuln.version)
378 else {
379 continue;
380 };
381 if manifests.insert((item.path.clone(), item.package.version.clone())) {
382 checks.push(still_locked_check(item.lockfile, &item.path, &item.package.name, &item.package.version));
383 }
384 if let Some(test) = test_command(item.lockfile, directory(&item.path)) {
385 tests.insert(test);
386 }
387 }
388 let locked: Vec<String> = manifests.iter().map(|(path, version)| format!("`{path}` ({version})")).collect();
389 body.push_str(&format!("\nLocked in: {}.\n", locked.join(", ")));
390 body.push_str(
391 "\nThe acceptance checks pass once no lockfile resolves a vulnerable version and the tests still pass. \
392 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\
393 ---\n_Opened by g1t's dependency upkeep. Turn it off for this project on its Security page._",
394 );
395 checks.extend(tests);
396 (body, checks)
397}
398
399#[cfg(test)]
400mod tests {
401 use super::*;
402 use g1t_contracts::security::LockfileText;
403
404 #[test]
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put405 fn scans_cost_their_cpu_and_the_rows_they_write() {
406 // 10 OSV calls and 3 vulnerabilities: 0.4 of CPU, 15 rows.
407 assert_eq!(dependency_check_cost(10, 3), 16);
408 assert_eq!(dependency_check_cost(0, 0), 2);
409 // A page of 25 commits that read 100 objects and found nothing:
410 // 10 of CPU and 2 rows. The old placeholder charged 100.
411 assert_eq!(crate::history::history_page_cost(100, 0), 12);
412 assert_eq!(crate::history::history_page_cost(0, 1), 3);
413 }
414
415 #[test]
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API416 fn go_sum_is_skipped_beside_go_mod() {
417 let files = Lockfiles {
418 commit: None,
419 files: vec![
420 LockfileText { path: "go.mod".into(), text: "require golang.org/x/net v0.7.0\n".into() },
421 LockfileText { path: "go.sum".into(), text: "golang.org/x/net v0.1.0 h1:x=\n".into() },
422 LockfileText { path: "tools/go.sum".into(), text: "golang.org/x/text v0.3.0 h1:x=\n".into() },
423 ],
424 };
425 let found: Vec<String> = packages(&files).iter().map(|l| format!("{}:{}", l.path, l.package.version)).collect();
426 assert_eq!(found, ["go.mod:v0.7.0", "tools/go.sum:v0.3.0"]);
427 }
428
429 #[test]
430 fn the_issue_names_the_advisories_and_checks_the_lockfile() {
431 let row = VulnRow {
432 id: "vul_1".into(),
433 repo_id: "rep_1".into(),
434 ecosystem: "npm".into(),
435 package: "lodash".into(),
436 version: "4.17.20".into(),
437 manifest: "web/package-lock.json".into(),
438 osv_id: "GHSA-35jh-r3h4-6jhm".into(),
439 advisory: "GHSA-35jh-r3h4-6jhm".into(),
440 summary: "Command Injection in lodash".into(),
441 severity: "high".into(),
442 fixed_version: Some("4.17.21".into()),
443 status: "open".into(),
444 found_at: "2026-10-04T00:00:00Z".into(),
445 fixed_at: None,
446 number: None,
447 };
448 let located = vec![Located {
449 package: Package { ecosystem: g1t_scan::lockfiles::Ecosystem::Npm, name: "lodash".into(), version: "4.17.20".into() },
450 lockfile: Lockfile::PackageLock,
451 path: "web/package-lock.json".into(),
452 }];
453 let (body, checks) = issue_text("npm", "lodash", "4.17.21", &[&row], &located);
454 assert!(body.contains("[GHSA-35jh-r3h4-6jhm](https://osv.dev/vulnerability/GHSA-35jh-r3h4-6jhm) | high | 4.17.20 | 4.17.21"));
455 assert!(body.contains("`web/package-lock.json` (4.17.20)"));
456 assert_eq!(checks.len(), 2);
457 assert!(checks[0].contains("node_modules/lodash") && checks[0].contains("'web/package-lock.json'"));
458 assert_eq!(checks[1], "cd 'web' && npm ci && npm test --if-present");
459 }
460}