Skip to content
789 linesCodeBlameRaw
1//! The security service: what g1t finds wrong in a repository, and the
2//! upkeep that fixes it without a person.
3//!
4//! - Secrets. The repos service refuses pushes that add one
5//! (`push_blocked` says which were allowed and records the rest), and
6//! each repository's history is scanned once, in the background.
7//! A push too large to scan first is let through, and its new commits
8//! are scanned after they land (`history::advance_push_scan`).
9//! - Alerts are open, dismissed (with a reason, a comment and who) or
10//! fixed, and each keeps an activity log. Likely test values are listed
11//! apart and never block a push or count as critical.
12//! - Dependencies. On every push to a default branch, and daily, the
13//! lockfiles are read and every package checked against OSV. Each
14//! vulnerable package with a fix gets a security update: g1t itself
15//! (`User::system`) opens a pull request raising its version, made in a
16//! sandbox (the runner's `bump`), which lands through the branch's
17//! required checks. Only when code has to change is g1t put on an
18//! issue for it.
19//!
20//! - The security suite (`g1t_contracts::security_suite`): custom secret
21//! patterns (`patterns`), push protection bypasses, their review and
22//! validity checks (`secret_alerts`), code scanning from SARIF uploads
23//! and its pull request check (`code_scanning`), the dependency graph,
24//! its SBOM and dependency review (`supply_chain`), "Fix with g1t"
25//! (`fixes`), and settings with the workspace's overview (`overview`).
26//! On private repositories its paid parts need the workspace's Security
27//! and quality activation (`suite`); it tells people through events.
28//!
29//! Other services reach it over `POST /rpc/<method>`; see
30//! `g1t_contracts::security`.
31
32mod code_scanning;
33mod config;
34mod deps;
35mod fixes;
36mod history;
37mod manifests;
38mod overview;
39mod patterns;
40mod planning;
41mod pull_text;
42mod ranges;
43mod registries;
44mod schedule;
45mod secret_alerts;
46mod security_updates;
47mod store;
48mod suite;
49mod suite_store;
50mod supply_chain;
51mod timezones;
52mod update_store;
53mod updates;
54mod version_updates;
55mod yaml;
56
57use g1t_contracts::access::{self, Capability};
58use g1t_contracts::events::{Event, WorkspaceRenamed};
59use g1t_contracts::security::UPDATE_BRANCH_PREFIX;
60use g1t_contracts::repos::{GetArgs, PathByIdArgs, Repo, RepoPath, RepoStatus, StatusByIdArgs};
61use g1t_contracts::security::*;
62use g1t_contracts::time::rfc3339;
63use g1t_contracts::{FailureCode, Outcome, User};
64use g1t_kit::{args, now_ms, reply, rpc_method};
65use serde::Deserialize;
66use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
67
68use store::{Activity, RepoRow, Store};
69
70/// Repositories whose history is continued per sweep, and pages each.
71const HISTORIES_PER_SWEEP: u32 = 5;
72const PAGES_PER_SWEEP: u32 = 4;
73/// Repositories whose dependencies are read again per sweep.
74const DEPENDENCIES_PER_SWEEP: u32 = 10;
75const DAY_MS: u64 = 24 * 60 * 60 * 1000;
76const MAX_REASON_CHARS: usize = 500;
77/// What seeing a repository's findings takes: the Write role, as changing
78/// its code does. Dismissing or allowing one takes Admin.
79const SEE_FINDINGS: Capability = Capability::Push;
80
81pub struct Security {
82 store: Store,
83 identity: Fetcher,
84 repos: Fetcher,
85 work: Fetcher,
86 runner: Fetcher,
87 billing: Fetcher,
88 actions: Fetcher,
89 /// The events service: security events for webhooks and the inbox,
90 /// and audit entries. Optional so a deployment without the binding
91 /// still scans.
92 events: Option<Fetcher>,
93}
94
95fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
96 Outcome::fail(code, message)
97}
98
99/// `git.push`, as far as this service reads it.
100#[derive(Deserialize)]
101#[serde(rename_all = "camelCase")]
102struct Pushed {
103 repo_id: String,
104 #[serde(default)]
105 default_branch: bool,
106 #[serde(default, rename = "ref")]
107 git_ref: String,
108 #[serde(default)]
109 before: Option<String>,
110 #[serde(default)]
111 after: String,
112 /// Too large to scan before it was stored.
113 #[serde(default)]
114 unscanned: bool,
115}
116
117/// `pull.merged`, `pull.closed` and `checks.completed`, as far as this
118/// service reads them.
119#[derive(Deserialize)]
120#[serde(rename_all = "camelCase")]
121struct PullHappened {
122 repo_id: String,
123 number: u32,
124 #[serde(default)]
125 status: Option<String>,
126}
127
128/// The longest comment a dismissal keeps.
129const MAX_COMMENT_CHARS: usize = MAX_REASON_CHARS;
130/// Activity rows the Security page reads, newest first.
131const ACTIVITY_SHOWN: u32 = 500;
132
133#[derive(Deserialize)]
134#[serde(rename_all = "camelCase")]
135struct Created {
136 repo_id: String,
137}
138
139impl Security {
140 fn new(env: &Env) -> Result<Self> {
141 Ok(Security {
142 store: Store { db: env.d1("DB")? },
143 identity: env.service("IDENTITY")?,
144 repos: env.service("REPOS")?,
145 work: env.service("WORK")?,
146 runner: env.service("RUNNER")?,
147 billing: env.service("BILLING")?,
148 actions: env.service("ACTIONS")?,
149 events: env.service("EVENTS").ok(),
150 })
151 }
152
153 /// The repository at `path`, recorded here, if `viewer` may see its
154 /// findings (the Write role) and do `capability`. Findings are for those
155 /// who can change the code: to anyone else the page does not exist,
156 /// public repository or not.
157 async fn member_repo(
158 &self,
159 path: &RepoPath,
160 viewer: &Option<User>,
161 capability: Capability,
162 ) -> Result<Outcome<RepoRow>> {
163 let hidden = || fail(FailureCode::NotFound, "Repository not found.");
164 if viewer.is_none() {
165 return Ok(hidden());
166 }
167 let repo: Outcome<Repo> =
168 g1t_kit::call(&self.repos, "get", &GetArgs { path: path.clone(), viewer: viewer.clone() }).await?;
169 let repo = match repo {
170 Outcome::Ok(repo) if repo.fork_of.is_none() => repo,
171 _ => return Ok(hidden()),
172 };
173 if !access::can(viewer.as_ref(), &repo, SEE_FINDINGS) {
174 return Ok(hidden());
175 }
176 if !access::can(viewer.as_ref(), &repo, capability) {
177 return Ok(fail(
178 FailureCode::Forbidden,
179 access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)),
180 ));
181 }
182 let row = self.store.register(&repo.id, &repo.namespace, &repo.name).await?;
183 self.store.set_private(&repo.id, repo.is_private).await?;
184 Ok(Outcome::Ok(row))
185 }
186
187 /// Whether the repository is neither archived nor deleted. When repos
188 /// cannot say, it is taken as active.
189 async fn active(&self, repo_id: &str) -> Result<bool> {
190 let status: Result<RepoStatus> =
191 g1t_kit::call(&self.repos, "status_by_id", &StatusByIdArgs { id: repo_id.to_owned() }).await;
192 Ok(match status {
193 Ok(status) => status.active(),
194 Err(error) => {
195 worker::console_error!("security: status_by_id {repo_id}: {error}");
196 true
197 }
198 })
199 }
200
201 /// Records a repository named in an event, by id. Forks are not
202 /// recorded: a pull request's findings belong to its repository.
203 async fn register_by_id(&self, repo_id: &str) -> Result<Option<RepoRow>> {
204 if let Some(row) = self.store.repo(repo_id).await? {
205 return Ok(Some(row));
206 }
207 let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await?;
208 match path {
209 Some(path) => Ok(Some(self.store.register(repo_id, &path.namespace, &path.name).await?)),
210 None => Ok(None),
211 }
212 }
213
214 async fn overview(&self, a: OverviewArgs) -> Result<Outcome<SecurityOverview>> {
215 let mut repo = match self.member_repo(&a.repo, &a.viewer, SEE_FINDINGS).await? {
216 Outcome::Ok(repo) => repo,
217 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
218 };
219 // The first look at a repository reads its dependencies at once;
220 // its history is scanned in the background.
221 if repo.deps_scanned_at.is_none() {
222 self.scan_dependencies(&repo).await?;
223 repo = self.store.repo(&repo.repo_id).await?.unwrap_or(repo);
224 }
225 let (counts, _, _) = self.store.counts(&repo.repo_id).await?;
226 Ok(Outcome::Ok(SecurityOverview {
227 repo_id: repo.repo_id.clone(),
228 counts,
229 secret_counts: self.store.secret_counts(&repo.repo_id).await?,
230 secrets: self.store.secrets(&repo.repo_id).await?,
231 vulnerabilities: self.store.vulnerabilities(&repo.repo_id).await?,
232 activity: self.store.activity(&repo.repo_id, ACTIVITY_SHOWN).await?,
233 scan: repo.scan_state(),
234 upkeep: repo.upkeep != 0,
235 version_updates: self.version_updates_view(&repo).await?,
236 }))
237 }
238
239 /// What dismissing or reopening alert `id` takes: Admin for a secret,
240 /// whose dismissal lets it through push protection, Write for a
241 /// vulnerable dependency.
242 fn capability_for(id: &str) -> Capability {
243 if id.starts_with("sec_") { Capability::ManageIntegrations } else { SEE_FINDINGS }
244 }
245
246 async fn dismiss(&self, a: DismissArgs) -> Result<Outcome<AlertChange>> {
247 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Self::capability_for(&a.id)).await? {
248 Outcome::Ok(repo) => repo,
249 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
250 };
251 if !a.actor.verified {
252 return Ok(fail(FailureCode::Forbidden, "Confirm your email address first."));
253 }
254 let comment: String = a.comment.trim().chars().take(MAX_COMMENT_CHARS).collect();
255 let comment = (!comment.is_empty()).then_some(comment);
256 if let Some(finding) = self.store.secret(&repo.repo_id, &a.id).await? {
257 if !a.reason.for_secrets() {
258 return Ok(fail(
259 FailureCode::Invalid,
260 "A secret is dismissed as false_positive, used_in_tests, revoked or wont_fix.",
261 ));
262 }
263 if finding.state != AlertState::Open {
264 return Ok(fail(FailureCode::Conflict, "This alert is not open. Reopen it first to dismiss it again."));
265 }
266 self.store
267 .dismiss_secret(&repo.repo_id, &a.id, a.reason, &a.actor.username, comment.as_deref())
268 .await?;
269 self.store
270 .record(&repo.repo_id, &[Activity {
271 alert_id: &a.id,
272 action: "dismissed",
273 actor: Some(&a.actor.username),
274 reason: Some(a.reason),
275 comment: comment.as_deref(),
276 number: None,
277 }])
278 .await?;
279 let secret = self.store.secret(&repo.repo_id, &a.id).await?;
280 if let Some(secret) = &secret {
281 // Revoked is fixed; any other reason, dismissed.
282 let action = if a.reason == DismissReason::Revoked { "fixed" } else { "dismissed" };
283 let event = g1t_contracts::security_suite::SecurityEvent {
284 reason: Some(a.reason.as_str().to_owned()),
285 ..secret_alerts::secret_event(&repo, secret)
286 };
287 self.alert_event(g1t_contracts::security_suite::AlertType::SecretScanning, action, &repo, event, Some(a.actor.id.clone()))
288 .await;
289 }
290 return Ok(Outcome::Ok(AlertChange { secret, vulnerability: None }));
291 }
292 let Some(vuln) = self.store.vulnerability(&repo.repo_id, &a.id).await? else {
293 return Ok(fail(FailureCode::NotFound, "No such alert."));
294 };
295 if a.reason.for_secrets() {
296 return Ok(fail(
297 FailureCode::Invalid,
298 "A dependency is dismissed as fix_started, no_bandwidth, tolerable_risk, inaccurate or not_used.",
299 ));
300 }
301 if vuln.state != AlertState::Open {
302 return Ok(fail(FailureCode::Conflict, "This alert is not open. Reopen it first to dismiss it again."));
303 }
304 self.store
305 .dismiss_vulnerability(&repo.repo_id, &a.id, a.reason, &a.actor.username, comment.as_deref())
306 .await?;
307 self.store
308 .record(&repo.repo_id, &[Activity {
309 alert_id: &a.id,
310 action: "dismissed",
311 actor: Some(&a.actor.username),
312 reason: Some(a.reason),
313 comment: comment.as_deref(),
314 number: None,
315 }])
316 .await?;
317 let vulnerability = self.store.vulnerability(&repo.repo_id, &a.id).await?;
318 if let Some(vuln) = &vulnerability {
319 let event = g1t_contracts::security_suite::SecurityEvent { reason: Some(a.reason.as_str().to_owned()), ..deps::vulnerability_event(&repo, vuln) };
320 self.alert_event(g1t_contracts::security_suite::AlertType::Vulnerability, "dismissed", &repo, event, Some(a.actor.id.clone())).await;
321 }
322 Ok(Outcome::Ok(AlertChange { secret: None, vulnerability }))
323 }
324
325 async fn reopen(&self, a: ReopenArgs) -> Result<Outcome<AlertChange>> {
326 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Self::capability_for(&a.id)).await? {
327 Outcome::Ok(repo) => repo,
328 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
329 };
330 if !a.actor.verified {
331 return Ok(fail(FailureCode::Forbidden, "Confirm your email address first."));
332 }
333 let reopened = Activity { alert_id: &a.id, action: "reopened", actor: Some(&a.actor.username), reason: None, comment: None, number: None };
334 if let Some(finding) = self.store.secret(&repo.repo_id, &a.id).await? {
335 if finding.state == AlertState::Open {
336 return Ok(fail(FailureCode::Conflict, "This alert is already open."));
337 }
338 // A secret that never landed has nothing to reopen to but blocked.
339 let status = if finding.source == "push" && finding.test_value.is_none() { SecretStatus::Blocked } else { SecretStatus::Open };
340 self.store.reopen_secret(&repo.repo_id, &a.id, status).await?;
341 self.store.record(&repo.repo_id, &[reopened]).await?;
342 let secret = self.store.secret(&repo.repo_id, &a.id).await?;
343 if let Some(secret) = &secret {
344 self.alert_event(
345 g1t_contracts::security_suite::AlertType::SecretScanning,
346 "reopened",
347 &repo,
348 secret_alerts::secret_event(&repo, secret),
349 Some(a.actor.id.clone()),
350 )
351 .await;
352 }
353 return Ok(Outcome::Ok(AlertChange { secret, vulnerability: None }));
354 }
355 let Some(vuln) = self.store.vulnerability(&repo.repo_id, &a.id).await? else {
356 return Ok(fail(FailureCode::NotFound, "No such alert."));
357 };
358 if vuln.state != AlertState::Dismissed {
359 return Ok(fail(FailureCode::Conflict, "Only a dismissed alert can be reopened; a fixed one reopens when it is found again."));
360 }
361 self.store.reopen_vulnerability(&repo.repo_id, &a.id).await?;
362 self.store.record(&repo.repo_id, &[reopened]).await?;
363 let vulnerability = self.store.vulnerability(&repo.repo_id, &a.id).await?;
364 if let Some(vuln) = &vulnerability {
365 self.alert_event(
366 g1t_contracts::security_suite::AlertType::Vulnerability,
367 "reopened",
368 &repo,
369 deps::vulnerability_event(&repo, vuln),
370 Some(a.actor.id.clone()),
371 )
372 .await;
373 }
374 Ok(Outcome::Ok(AlertChange { secret: None, vulnerability }))
375 }
376
377 async fn rescan(&self, a: RescanArgs) -> Result<Outcome<ScanState>> {
378 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), SEE_FINDINGS).await? {
379 Outcome::Ok(repo) => repo,
380 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
381 };
382 self.store.restart_history(&repo.repo_id).await?;
383 self.scan_dependencies(&repo).await?;
384 if let Some(fresh) = self.store.repo(&repo.repo_id).await? {
385 self.advance_history(&fresh, 1).await?;
386 }
387 let state = self.store.repo(&repo.repo_id).await?.map(|row| row.scan_state()).unwrap_or_default();
388 Ok(Outcome::Ok(state))
389 }
390
391 async fn set_upkeep(&self, a: SetUpkeepArgs) -> Result<Outcome<bool>> {
392 // Whether agents keep its dependencies up to date is one of its
393 // settings.
394 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Capability::ManageSettings).await? {
395 Outcome::Ok(repo) => repo,
396 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
397 };
398 if !a.actor.verified {
399 return Ok(fail(FailureCode::Forbidden, "Confirm your email address first."));
400 }
401 self.store.set_upkeep(&repo.repo_id, a.enabled, &a.actor.username).await?;
402 Ok(Outcome::Ok(a.enabled))
403 }
404
405 async fn workspace(&self, a: WorkspaceArgs) -> Result<Outcome<Vec<RepoSecurity>>> {
406 let workspace = a.workspace.to_lowercase();
407 if !a.viewer.as_ref().is_some_and(|user| user.is_member(&workspace)) {
408 return Ok(fail(FailureCode::NotFound, "Workspace not found."));
409 }
410 let mut list = Vec::new();
411 for repo in self.store.in_namespace(&workspace).await? {
412 // Only the repositories whose findings the viewer may see. The
413 // Write role is never had through being public, so treating
414 // each as private changes nothing.
415 let target = access::RepoRef { id: &repo.repo_id, namespace: &workspace, private: true };
416 if !access::can(a.viewer.as_ref(), target, SEE_FINDINGS) {
417 continue;
418 }
419 let (counts, secrets, vulnerabilities) = self.store.counts(&repo.repo_id).await?;
420 list.push(RepoSecurity {
421 repo_id: repo.repo_id,
422 name: repo.name,
423 counts,
424 secrets,
425 vulnerabilities,
426 upkeep: repo.upkeep != 0,
427 dependencies_scanned_at: repo.deps_scanned_at,
428 });
429 }
430 Ok(Outcome::Ok(list))
431 }
432
433 /// Push protection's question: which of these secrets were allowed?
434 /// The others are recorded as blocked, so someone can allow them.
435 async fn push_blocked(&self, a: PushBlockedArgs) -> Result<PushVerdict> {
436 let repo = self.store.register(&a.repo_id, &a.path.namespace, &a.path.name).await?;
437 if let Some(private) = a.private {
438 self.store.set_private(&a.repo_id, private).await?;
439 }
440 let fingerprints: Vec<String> = a.secrets.iter().map(|secret| secret.fingerprint.clone()).collect();
441 let known = self.store.known(&a.repo_id, &fingerprints).await?;
442 let mut allowed = let_through(&known, &a.secrets);
443 // Bypassed with a reason: let through, whatever the alert says now.
444 allowed.extend(self.store.bypassed(&a.repo_id, &fingerprints).await?);
445 allowed.sort();
446 allowed.dedup();
447 let all = a.secrets.clone();
448 let fresh: Vec<NewSecret> = a
449 .secrets
450 .into_iter()
451 .filter(|secret| !known.iter().any(|(fingerprint, _, _)| *fingerprint == secret.fingerprint))
452 .collect();
453 // A likely test value goes through, so it lands: open, not blocked.
454 let (tests, real): (Vec<NewSecret>, Vec<NewSecret>) = fresh.into_iter().partition(|secret| secret.test_value.is_some());
455 self.store
456 .add_secrets(&a.repo_id, &real, SecretStatus::Blocked, "push", a.pusher.as_deref())
457 .await?;
458 self.store
459 .add_secrets(&a.repo_id, &tests, SecretStatus::Open, "push", a.pusher.as_deref())
460 .await?;
461 let fresh: Vec<String> = real.iter().map(|secret| secret.fingerprint.clone()).collect();
462 self.secrets_found(&repo, &all, "push", &fresh, a.pusher.as_deref()).await?;
463 let ids = self
464 .store
465 .known(&a.repo_id, &fingerprints)
466 .await?
467 .into_iter()
468 .filter(|(fingerprint, _, _)| !allowed.contains(fingerprint))
469 .map(|(fingerprint, id, _)| (fingerprint, id))
470 .collect();
471 Ok(PushVerdict { allowed, ids })
472 }
473
474 /// What happens on the bus that concerns this service.
475 async fn on_event(&self, event: &Event) -> Result<()> {
476 match event.kind.as_str() {
477 "git.push" => {
478 let Ok(pushed) = serde_json::from_value::<Pushed>(event.data.clone()) else {
479 return Ok(());
480 };
481 // Too large to scan before it was stored: its new commits
482 // are scanned now, on whichever branch.
483 if pushed.unscanned
484 && !pushed.after.is_empty()
485 && let Some(repo) = self.register_by_id(&pushed.repo_id).await?
486 {
487 let id = self
488 .store
489 .add_push_scan(&repo.repo_id, &pushed.git_ref, &pushed.after, pushed.before.as_deref(), event.actor.as_deref())
490 .await?;
491 self.advance_push_scan(&id, history::PUSH_PAGES_AT_ONCE).await?;
492 }
493 // A version update's branch, or a grouped security update's,
494 // pushed by its sandbox: time for its pull request.
495 if !pushed.default_branch
496 && let Some(branch) = pushed.git_ref.strip_prefix("refs/heads/")
497 && self.update_pull_pushed(&pushed.repo_id, branch, &pushed.after).await?
498 {
499 return Ok(());
500 }
501 // A security update's branch, pushed by its sandbox: time for
502 // its pull request.
503 if let Some(branch) = pushed.git_ref.strip_prefix("refs/heads/")
504 && branch.starts_with(UPDATE_BRANCH_PREFIX)
505 {
506 self.update_pushed(&pushed.repo_id, branch).await?;
507 return Ok(());
508 }
509 if !pushed.default_branch {
510 return Ok(());
511 }
512 if let Some(repo) = self.register_by_id(&pushed.repo_id).await? {
513 self.scan_dependencies(&repo).await?;
514 if repo.history != "done" {
515 self.advance_history(&repo, 1).await?;
516 }
517 }
518 }
519 "pull.merged" | "pull.closed" | "checks.completed" => {
520 if let Ok(happened) = serde_json::from_value::<PullHappened>(event.data.clone())
521 && !(event.kind != "checks.completed" && self.update_pull_closed(&event.kind, &happened.repo_id, happened.number).await?)
522 {
523 self.update_pull_event(
524 &event.kind,
525 &happened.repo_id,
526 happened.number,
527 happened.status.as_deref(),
528 event.actor.as_deref(),
529 )
530 .await?;
531 }
532 }
533 "comment.created" => self.update_comment(event).await?,
534 // A pull request opened or its head moved: a dependency update
535 // file it changes is checked (not on `pull.ready`, which moves
536 // nothing), and dependency review runs.
537 "pull.opened" | "pull.updated" | "pull.ready" => {
538 if event.kind != "pull.ready" {
539 self.check_dependabot_file(event).await?;
540 }
541 if let Ok(happened) = serde_json::from_value::<PullHappened>(event.data.clone()) {
542 self.review_pull(&happened.repo_id, happened.number).await?;
543 }
544 }
545 "repo.created" => {
546 if let Ok(created) = serde_json::from_value::<Created>(event.data.clone()) {
547 self.register_by_id(&created.repo_id).await?;
548 }
549 }
550 "repo.visibility_changed" => {
551 if let Ok(changed) = serde_json::from_value::<g1t_contracts::events::RepoVisibilityChanged>(event.data.clone())
552 && self.store.repo(&changed.repo_id).await?.is_some()
553 {
554 self.store.set_private(&changed.repo_id, changed.is_private).await?;
555 }
556 }
557 // A repository transferred or renamed: it is recorded under its new path.
558 "repo.transferred" | "repo.renamed" => {
559 if let Some(moved) = g1t_kit::transfer::read(event) {
560 // Where it is now, so moves heard out of order end in
561 // the same place.
562 let now: Option<g1t_contracts::repos::RepoPath> = g1t_kit::call(
563 &self.repos,
564 "path_by_id",
565 &g1t_contracts::repos::PathByIdArgs { id: moved.repo_id.clone() },
566 )
567 .await?;
568 let current = now.map_or_else(
569 || moved.destination().to_owned(),
570 |path| format!("{}/{}", path.namespace, path.name),
571 );
572 if let Some((namespace, name)) = current.split_once('/') {
573 self.store.moved(&moved.repo_id, namespace, name).await?;
574 }
575 }
576 }
577 // A repository purged: everything found in it goes.
578 "repo.purged" => {
579 if let Some(g1t_kit::lifecycle::Lifecycle::Purged(purged)) = g1t_kit::lifecycle::read(event) {
580 self.store.purge_suite(&purged.repo_id).await?;
581 self.store.purge(&purged.repo_id).await?;
582 }
583 }
584 "workspace.renamed" => {
585 if let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) {
586 self.store.rename_namespace(&renamed.stale_slugs(&renamed.to), &renamed.to).await?;
587 }
588 }
589 _ => {}
590 }
591 Ok(())
592 }
593
594 /// The sweep: continues history scans and scans of pushes that landed
595 /// unscanned, catches security updates whose sandbox never pushed, and
596 /// reads dependencies that have not been read for a day.
597 async fn sweep(&self) -> Result<()> {
598 for scan in self.store.pending_push_scans(HISTORIES_PER_SWEEP).await? {
599 if let Err(error) = self.advance_push_scan(&scan.id, PAGES_PER_SWEEP).await {
600 worker::console_error!("security: push scan {} not continued: {error}", scan.id);
601 }
602 }
603 if let Err(error) = self.stalled_updates().await {
604 worker::console_error!("security: stalled security updates not handled: {error}");
605 }
606 if let Err(error) = self.sweep_suite().await {
607 worker::console_error!("security: daily snapshots and validity checks: {error}");
608 }
609 // Archived and deleted repositories wait; a few more are looked at
610 // so that they do not hold up the rest.
611 let mut histories = 0;
612 for repo in self.store.unfinished_histories(HISTORIES_PER_SWEEP * 4).await? {
613 if histories == HISTORIES_PER_SWEEP {
614 break;
615 }
616 if !self.active(&repo.repo_id).await? {
617 continue;
618 }
619 histories += 1;
620 if let Err(error) = self.advance_history(&repo, PAGES_PER_SWEEP).await {
621 worker::console_error!("security: history of {} not scanned: {error}", repo.repo_id);
622 }
623 }
624 let day_ago = rfc3339(now_ms().saturating_sub(DAY_MS));
625 for repo in self.store.stale_dependencies(&day_ago, DEPENDENCIES_PER_SWEEP).await? {
626 if !self.active(&repo.repo_id).await? {
627 self.store
628 .skip_dependencies(&repo.repo_id, "Dependencies are not checked while the repository is archived or deleted.")
629 .await?;
630 continue;
631 }
632 if let Err(error) = self.scan_dependencies(&repo).await {
633 worker::console_error!("security: dependencies of {} not read: {error}", repo.repo_id);
634 }
635 }
636 Ok(())
637 }
638}
639
640/// The fingerprints a push may carry: those someone dismissed (allowed),
641/// and likely test values, which are recorded but never stop a push.
642/// `known` is (fingerprint, id, status) of findings already recorded.
643fn let_through(known: &[(String, String, String)], secrets: &[NewSecret]) -> Vec<String> {
644 let mut allowed: Vec<String> = known
645 .iter()
646 .filter(|(_, _, status)| status == "allowed")
647 .map(|(fingerprint, _, _)| fingerprint.clone())
648 .chain(secrets.iter().filter(|secret| secret.test_value.is_some()).map(|secret| secret.fingerprint.clone()))
649 .collect();
650 allowed.sort();
651 allowed.dedup();
652 allowed
653}
654
655#[event(fetch)]
656async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
657 let Some(method) = rpc_method(&request) else {
658 return Response::error("Not found", 404);
659 };
660 let security = Security::new(&env)?;
661 let body: serde_json::Value = request.json().await?;
662 match method.as_str() {
663 "overview" => reply(&security.overview(args(body)?).await?),
664 "dismiss" => reply(&security.dismiss(args(body)?).await?),
665 "reopen" => reply(&security.reopen(args(body)?).await?),
666 "rescan" => reply(&security.rescan(args(body)?).await?),
667 "set_upkeep" => reply(&security.set_upkeep(args(body)?).await?),
668 "workspace" => reply(&security.workspace(args(body)?).await?),
669 "push_blocked" => reply(&security.push_blocked(args(body)?).await?),
670 "check_updates" => reply(&security.check_updates(args(body)?).await?),
671 // The security suite.
672 "patterns_for" => reply(&security.patterns_for(args(body)?).await?),
673 "custom_patterns" => reply(&security.custom_patterns(args(body)?).await?),
674 "save_custom_pattern" => reply(&security.save_custom_pattern(args(body)?).await?),
675 "delete_custom_pattern" => reply(&security.delete_custom_pattern(args(body)?).await?),
676 "dry_run_pattern" => reply(&security.dry_run_pattern(args(body)?).await?),
677 "secret_alert" => reply(&security.secret_alert(args(body)?).await?),
678 "bypass" => reply(&security.bypass(args(body)?).await?),
679 "bypass_requests" => reply(&security.bypass_requests(args(body)?).await?),
680 "review_bypass" => reply(&security.review_bypass(args(body)?).await?),
681 "check_validity" => reply(&security.check_validity(args(body)?).await?),
682 "upload_sarif" => reply(&security.upload_sarif(args(body)?).await?),
683 "sarif_status" => reply(&security.sarif_status(args(body)?).await?),
684 "code_scanning" => reply(&security.code_scanning(args(body)?).await?),
685 "code_alert" => reply(&security.code_alert(args(body)?).await?),
686 "set_code_alert_state" => reply(&security.set_code_alert_state(args(body)?).await?),
687 "pull_code_scanning" => reply(&security.pull_code_scanning(args(body)?).await?),
688 "fix_alert" => reply(&security.fix_alert(args(body)?).await?),
689 "dependency_graph" => reply(&security.dependency_graph(args(body)?).await?),
690 "sbom" => reply(&security.sbom(args(body)?).await?),
691 "dependency_review" => reply(&security.dependency_review(args(body)?).await?),
692 "security_settings" => reply(&security.security_settings(args(body)?).await?),
693 "set_security_settings" => reply(&security.set_security_settings(args(body)?).await?),
694 "workspace_security_settings" => reply(&security.workspace_security_settings(args(body)?).await?),
695 "set_workspace_security_settings" => reply(&security.set_workspace_security_settings(args(body)?).await?),
696 "security_overview" => reply(&security.security_overview(args(body)?).await?),
697 "workspace_alerts" => reply(&security.workspace_alerts(args(body)?).await?),
698 _ => Response::error("Unknown method", 404),
699 }
700}
701
702#[event(queue)]
703async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
704 let security = Security::new(&env)?;
705 for message in batch.messages()? {
706 let event = message.body();
707 if let Err(error) = security.on_event(event).await {
708 // Scans are idempotent and the sweep catches up, so one failed
709 // event is logged rather than retried.
710 worker::console_error!("security: {} {} failed: {error}", event.kind, event.id);
711 }
712 }
713 Ok(())
714}
715
716/// The cron (wrangler.jsonc) that runs version updates that are due.
717const VERSION_UPDATES_CRON: &str = "*/5 * * * *";
718
719#[event(scheduled)]
720async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
721 match Security::new(&env) {
722 // Version updates run every few minutes, so a schedule's time is kept.
723 Ok(security) if event.cron() == VERSION_UPDATES_CRON => {
724 if let Err(error) = security.version_update_sweep().await {
725 worker::console_error!("security: the version update sweep failed: {error}");
726 }
727 }
728 Ok(security) => {
729 if let Err(error) = security.sweep().await {
730 worker::console_error!("security: the sweep failed: {error}");
731 }
732 }
733 Err(error) => worker::console_error!("security: could not start: {error}"),
734 }
735}
736
737#[cfg(test)]
738mod tests {
739 use super::*;
740
741 fn secret(fingerprint: &str, test_value: Option<&str>) -> NewSecret {
742 NewSecret {
743 fingerprint: fingerprint.into(),
744 kind: "aws_access_key".into(),
745 path: "a.env".into(),
746 line: 1,
747 commit: "c".into(),
748 preview: "AKIA…".into(),
749 test_value: test_value.map(str::to_owned),
750 pattern_id: None,
751 pattern_name: None,
752 }
753 }
754
755 #[test]
756 fn dismissed_secrets_and_test_values_go_through() {
757 let known = vec![
758 ("allowed".to_owned(), "sec_1".to_owned(), "allowed".to_owned()),
759 ("resolved".to_owned(), "sec_2".to_owned(), "resolved".to_owned()),
760 ("blocked".to_owned(), "sec_3".to_owned(), "blocked".to_owned()),
761 ];
762 let secrets = [secret("allowed", None), secret("resolved", None), secret("example", Some("it says it is an example")), secret("real", None)];
763 assert_eq!(let_through(&known, &secrets), ["allowed", "example"]);
764 }
765
766 #[test]
767 fn a_push_says_when_it_landed_unscanned() {
768 let pushed: Pushed = serde_json::from_value(serde_json::json!({
769 "repoId": "rep_1", "ref": "refs/heads/import", "after": "abc", "defaultBranch": false, "unscanned": true
770 }))
771 .unwrap();
772 assert!(pushed.unscanned && pushed.before.is_none() && pushed.git_ref == "refs/heads/import");
773 let ordinary: Pushed = serde_json::from_value(serde_json::json!({ "repoId": "rep_1", "ref": "refs/heads/main", "after": "abc", "defaultBranch": true })).unwrap();
774 assert!(!ordinary.unscanned);
775 }
776
777 #[test]
778 fn owners_are_told_what_landed_and_what_to_do() {
779 let one = history::landed_secrets_intro("acme", "rocket", "import", 1);
780 assert!(one.contains("A push to import in acme/rocket") && one.contains("a secret that looks real") && one.contains("Rotate"));
781 assert!(history::landed_secrets_intro("acme", "rocket", "main", 3).contains("3 secrets that look real"));
782 }
783
784 #[test]
785 fn dismissing_a_secret_takes_admin_and_a_dependency_write() {
786 assert_eq!(Security::capability_for("sec_1"), Capability::ManageIntegrations);
787 assert_eq!(Security::capability_for("vul_1"), Capability::Push);
788 }
789}