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