g1t/services/security/src/lib.rs
| 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-agent put on an |
| 18 | //! issue for it. |
| 19 | //! |
| 20 | //! Other services reach it over `POST /rpc/<method>`; see |
| 21 | //! `g1t_contracts::security`. |
| 22 | |
| 23 | mod deps; |
| 24 | mod history; |
| 25 | mod security_updates; |
| 26 | mod store; |
| 27 | mod updates; |
| 28 | |
| 29 | use g1t_contracts::access::{self, Capability}; |
| 30 | use g1t_contracts::events::{Event, WorkspaceRenamed}; |
| 31 | use g1t_contracts::security::UPDATE_BRANCH_PREFIX; |
| 32 | use g1t_contracts::repos::{GetArgs, PathByIdArgs, Repo, RepoPath, RepoStatus, StatusByIdArgs}; |
| 33 | use g1t_contracts::security::*; |
| 34 | use g1t_contracts::time::rfc3339; |
| 35 | use g1t_contracts::{FailureCode, Outcome, User}; |
| 36 | use g1t_kit::{args, now_ms, reply, rpc_method}; |
| 37 | use serde::Deserialize; |
| 38 | use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event}; |
| 39 | |
| 40 | use store::{Activity, RepoRow, Store}; |
| 41 | |
| 42 | /// Repositories whose history is continued per sweep, and pages each. |
| 43 | const HISTORIES_PER_SWEEP: u32 = 5; |
| 44 | const PAGES_PER_SWEEP: u32 = 4; |
| 45 | /// Repositories whose dependencies are read again per sweep. |
| 46 | const DEPENDENCIES_PER_SWEEP: u32 = 10; |
| 47 | const DAY_MS: u64 = 24 * 60 * 60 * 1000; |
| 48 | const MAX_REASON_CHARS: usize = 500; |
| 49 | /// What seeing a repository's findings takes: the Write role, as changing |
| 50 | /// its code does. Dismissing or allowing one takes Admin. |
| 51 | const SEE_FINDINGS: Capability = Capability::Push; |
| 52 | |
| 53 | pub struct Security { |
| 54 | store: Store, |
| 55 | identity: Fetcher, |
| 56 | repos: Fetcher, |
| 57 | work: Fetcher, |
| 58 | runner: Fetcher, |
| 59 | billing: Fetcher, |
| 60 | } |
| 61 | |
| 62 | fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> { |
| 63 | Outcome::fail(code, message) |
| 64 | } |
| 65 | |
| 66 | /// `git.push`, as far as this service reads it. |
| 67 | #[derive(Deserialize)] |
| 68 | #[serde(rename_all = "camelCase")] |
| 69 | struct Pushed { |
| 70 | repo_id: String, |
| 71 | #[serde(default)] |
| 72 | default_branch: bool, |
| 73 | #[serde(default, rename = "ref")] |
| 74 | git_ref: String, |
| 75 | #[serde(default)] |
| 76 | before: Option<String>, |
| 77 | #[serde(default)] |
| 78 | after: String, |
| 79 | /// Too large to scan before it was stored. |
| 80 | #[serde(default)] |
| 81 | unscanned: bool, |
| 82 | } |
| 83 | |
| 84 | /// `pull.merged`, `pull.closed` and `checks.completed`, as far as this |
| 85 | /// service reads them. |
| 86 | #[derive(Deserialize)] |
| 87 | #[serde(rename_all = "camelCase")] |
| 88 | struct PullHappened { |
| 89 | repo_id: String, |
| 90 | number: u32, |
| 91 | #[serde(default)] |
| 92 | status: Option<String>, |
| 93 | } |
| 94 | |
| 95 | /// The longest comment a dismissal keeps. |
| 96 | const MAX_COMMENT_CHARS: usize = MAX_REASON_CHARS; |
| 97 | /// Activity rows the Security page reads, newest first. |
| 98 | const ACTIVITY_SHOWN: u32 = 500; |
| 99 | |
| 100 | #[derive(Deserialize)] |
| 101 | #[serde(rename_all = "camelCase")] |
| 102 | struct Created { |
| 103 | repo_id: String, |
| 104 | } |
| 105 | |
| 106 | impl Security { |
| 107 | fn new(env: &Env) -> Result<Self> { |
| 108 | Ok(Security { |
| 109 | store: Store { db: env.d1("DB")? }, |
| 110 | identity: env.service("IDENTITY")?, |
| 111 | repos: env.service("REPOS")?, |
| 112 | work: env.service("WORK")?, |
| 113 | runner: env.service("RUNNER")?, |
| 114 | billing: env.service("BILLING")?, |
| 115 | }) |
| 116 | } |
| 117 | |
| 118 | /// The repository at `path`, recorded here, if `viewer` may see its |
| 119 | /// findings (the Write role) and do `capability`. Findings are for those |
| 120 | /// who can change the code: to anyone else the page does not exist, |
| 121 | /// public repository or not. |
| 122 | async fn member_repo( |
| 123 | &self, |
| 124 | path: &RepoPath, |
| 125 | viewer: &Option<User>, |
| 126 | capability: Capability, |
| 127 | ) -> Result<Outcome<RepoRow>> { |
| 128 | let hidden = || fail(FailureCode::NotFound, "Repository not found."); |
| 129 | if viewer.is_none() { |
| 130 | return Ok(hidden()); |
| 131 | } |
| 132 | let repo: Outcome<Repo> = |
| 133 | g1t_kit::call(&self.repos, "get", &GetArgs { path: path.clone(), viewer: viewer.clone() }).await?; |
| 134 | let repo = match repo { |
| 135 | Outcome::Ok(repo) if repo.fork_of.is_none() => repo, |
| 136 | _ => return Ok(hidden()), |
| 137 | }; |
| 138 | if !access::can(viewer.as_ref(), &repo, SEE_FINDINGS) { |
| 139 | return Ok(hidden()); |
| 140 | } |
| 141 | if !access::can(viewer.as_ref(), &repo, capability) { |
| 142 | return Ok(fail( |
| 143 | FailureCode::Forbidden, |
| 144 | access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)), |
| 145 | )); |
| 146 | } |
| 147 | Ok(Outcome::Ok(self.store.register(&repo.id, &repo.namespace, &repo.name).await?)) |
| 148 | } |
| 149 | |
| 150 | /// Whether the repository is neither archived nor deleted. When repos |
| 151 | /// cannot say, it is taken as active. |
| 152 | async fn active(&self, repo_id: &str) -> Result<bool> { |
| 153 | let status: Result<RepoStatus> = |
| 154 | g1t_kit::call(&self.repos, "status_by_id", &StatusByIdArgs { id: repo_id.to_owned() }).await; |
| 155 | Ok(match status { |
| 156 | Ok(status) => status.active(), |
| 157 | Err(error) => { |
| 158 | worker::console_error!("security: status_by_id {repo_id}: {error}"); |
| 159 | true |
| 160 | } |
| 161 | }) |
| 162 | } |
| 163 | |
| 164 | /// Records a repository named in an event, by id. Forks are not |
| 165 | /// recorded: a pull request's findings belong to its repository. |
| 166 | async fn register_by_id(&self, repo_id: &str) -> Result<Option<RepoRow>> { |
| 167 | if let Some(row) = self.store.repo(repo_id).await? { |
| 168 | return Ok(Some(row)); |
| 169 | } |
| 170 | let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await?; |
| 171 | match path { |
| 172 | Some(path) => Ok(Some(self.store.register(repo_id, &path.namespace, &path.name).await?)), |
| 173 | None => Ok(None), |
| 174 | } |
| 175 | } |
| 176 | |
| 177 | async fn overview(&self, a: OverviewArgs) -> Result<Outcome<SecurityOverview>> { |
| 178 | let mut repo = match self.member_repo(&a.repo, &a.viewer, SEE_FINDINGS).await? { |
| 179 | Outcome::Ok(repo) => repo, |
| 180 | Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)), |
| 181 | }; |
| 182 | // The first look at a repository reads its dependencies at once; |
| 183 | // its history is scanned in the background. |
| 184 | if repo.deps_scanned_at.is_none() { |
| 185 | self.scan_dependencies(&repo).await?; |
| 186 | repo = self.store.repo(&repo.repo_id).await?.unwrap_or(repo); |
| 187 | } |
| 188 | let (counts, _, _) = self.store.counts(&repo.repo_id).await?; |
| 189 | Ok(Outcome::Ok(SecurityOverview { |
| 190 | repo_id: repo.repo_id.clone(), |
| 191 | counts, |
| 192 | secret_counts: self.store.secret_counts(&repo.repo_id).await?, |
| 193 | secrets: self.store.secrets(&repo.repo_id).await?, |
| 194 | vulnerabilities: self.store.vulnerabilities(&repo.repo_id).await?, |
| 195 | activity: self.store.activity(&repo.repo_id, ACTIVITY_SHOWN).await?, |
| 196 | scan: repo.scan_state(), |
| 197 | upkeep: repo.upkeep != 0, |
| 198 | version_updates: repo.version_updates(), |
| 199 | })) |
| 200 | } |
| 201 | |
| 202 | /// What dismissing or reopening alert `id` takes: Admin for a secret, |
| 203 | /// whose dismissal lets it through push protection, Write for a |
| 204 | /// vulnerable dependency. |
| 205 | fn capability_for(id: &str) -> Capability { |
| 206 | if id.starts_with("sec_") { Capability::ManageIntegrations } else { SEE_FINDINGS } |
| 207 | } |
| 208 | |
| 209 | async fn dismiss(&self, a: DismissArgs) -> Result<Outcome<AlertChange>> { |
| 210 | let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Self::capability_for(&a.id)).await? { |
| 211 | Outcome::Ok(repo) => repo, |
| 212 | Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)), |
| 213 | }; |
| 214 | if !a.actor.verified { |
| 215 | return Ok(fail(FailureCode::Forbidden, "Confirm your email address first.")); |
| 216 | } |
| 217 | let comment: String = a.comment.trim().chars().take(MAX_COMMENT_CHARS).collect(); |
| 218 | let comment = (!comment.is_empty()).then_some(comment); |
| 219 | if let Some(finding) = self.store.secret(&repo.repo_id, &a.id).await? { |
| 220 | if !a.reason.for_secrets() { |
| 221 | return Ok(fail( |
| 222 | FailureCode::Invalid, |
| 223 | "A secret is dismissed as false_positive, used_in_tests, revoked or wont_fix.", |
| 224 | )); |
| 225 | } |
| 226 | if finding.state != AlertState::Open { |
| 227 | return Ok(fail(FailureCode::Conflict, "This alert is not open. Reopen it first to dismiss it again.")); |
| 228 | } |
| 229 | self.store |
| 230 | .dismiss_secret(&repo.repo_id, &a.id, a.reason, &a.actor.username, comment.as_deref()) |
| 231 | .await?; |
| 232 | self.store |
| 233 | .record(&repo.repo_id, &[Activity { |
| 234 | alert_id: &a.id, |
| 235 | action: "dismissed", |
| 236 | actor: Some(&a.actor.username), |
| 237 | reason: Some(a.reason), |
| 238 | comment: comment.as_deref(), |
| 239 | number: None, |
| 240 | }]) |
| 241 | .await?; |
| 242 | return Ok(Outcome::Ok(AlertChange { secret: self.store.secret(&repo.repo_id, &a.id).await?, vulnerability: None })); |
| 243 | } |
| 244 | let Some(vuln) = self.store.vulnerability(&repo.repo_id, &a.id).await? else { |
| 245 | return Ok(fail(FailureCode::NotFound, "No such alert.")); |
| 246 | }; |
| 247 | if a.reason.for_secrets() { |
| 248 | return Ok(fail( |
| 249 | FailureCode::Invalid, |
| 250 | "A dependency is dismissed as fix_started, no_bandwidth, tolerable_risk, inaccurate or not_used.", |
| 251 | )); |
| 252 | } |
| 253 | if vuln.state != AlertState::Open { |
| 254 | return Ok(fail(FailureCode::Conflict, "This alert is not open. Reopen it first to dismiss it again.")); |
| 255 | } |
| 256 | self.store |
| 257 | .dismiss_vulnerability(&repo.repo_id, &a.id, a.reason, &a.actor.username, comment.as_deref()) |
| 258 | .await?; |
| 259 | self.store |
| 260 | .record(&repo.repo_id, &[Activity { |
| 261 | alert_id: &a.id, |
| 262 | action: "dismissed", |
| 263 | actor: Some(&a.actor.username), |
| 264 | reason: Some(a.reason), |
| 265 | comment: comment.as_deref(), |
| 266 | number: None, |
| 267 | }]) |
| 268 | .await?; |
| 269 | Ok(Outcome::Ok(AlertChange { secret: None, vulnerability: self.store.vulnerability(&repo.repo_id, &a.id).await? })) |
| 270 | } |
| 271 | |
| 272 | async fn reopen(&self, a: ReopenArgs) -> Result<Outcome<AlertChange>> { |
| 273 | let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Self::capability_for(&a.id)).await? { |
| 274 | Outcome::Ok(repo) => repo, |
| 275 | Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)), |
| 276 | }; |
| 277 | if !a.actor.verified { |
| 278 | return Ok(fail(FailureCode::Forbidden, "Confirm your email address first.")); |
| 279 | } |
| 280 | let reopened = Activity { alert_id: &a.id, action: "reopened", actor: Some(&a.actor.username), reason: None, comment: None, number: None }; |
| 281 | if let Some(finding) = self.store.secret(&repo.repo_id, &a.id).await? { |
| 282 | if finding.state == AlertState::Open { |
| 283 | return Ok(fail(FailureCode::Conflict, "This alert is already open.")); |
| 284 | } |
| 285 | // A secret that never landed has nothing to reopen to but blocked. |
| 286 | let status = if finding.source == "push" && finding.test_value.is_none() { SecretStatus::Blocked } else { SecretStatus::Open }; |
| 287 | self.store.reopen_secret(&repo.repo_id, &a.id, status).await?; |
| 288 | self.store.record(&repo.repo_id, &[reopened]).await?; |
| 289 | return Ok(Outcome::Ok(AlertChange { secret: self.store.secret(&repo.repo_id, &a.id).await?, vulnerability: None })); |
| 290 | } |
| 291 | let Some(vuln) = self.store.vulnerability(&repo.repo_id, &a.id).await? else { |
| 292 | return Ok(fail(FailureCode::NotFound, "No such alert.")); |
| 293 | }; |
| 294 | if vuln.state != AlertState::Dismissed { |
| 295 | return Ok(fail(FailureCode::Conflict, "Only a dismissed alert can be reopened; a fixed one reopens when it is found again.")); |
| 296 | } |
| 297 | self.store.reopen_vulnerability(&repo.repo_id, &a.id).await?; |
| 298 | self.store.record(&repo.repo_id, &[reopened]).await?; |
| 299 | Ok(Outcome::Ok(AlertChange { secret: None, vulnerability: self.store.vulnerability(&repo.repo_id, &a.id).await? })) |
| 300 | } |
| 301 | |
| 302 | async fn rescan(&self, a: RescanArgs) -> Result<Outcome<ScanState>> { |
| 303 | let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), SEE_FINDINGS).await? { |
| 304 | Outcome::Ok(repo) => repo, |
| 305 | Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)), |
| 306 | }; |
| 307 | self.store.restart_history(&repo.repo_id).await?; |
| 308 | self.scan_dependencies(&repo).await?; |
| 309 | if let Some(fresh) = self.store.repo(&repo.repo_id).await? { |
| 310 | self.advance_history(&fresh, 1).await?; |
| 311 | } |
| 312 | let state = self.store.repo(&repo.repo_id).await?.map(|row| row.scan_state()).unwrap_or_default(); |
| 313 | Ok(Outcome::Ok(state)) |
| 314 | } |
| 315 | |
| 316 | async fn set_upkeep(&self, a: SetUpkeepArgs) -> Result<Outcome<bool>> { |
| 317 | // Whether agents keep its dependencies up to date is one of its |
| 318 | // settings. |
| 319 | let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Capability::ManageSettings).await? { |
| 320 | Outcome::Ok(repo) => repo, |
| 321 | Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)), |
| 322 | }; |
| 323 | if !a.actor.verified { |
| 324 | return Ok(fail(FailureCode::Forbidden, "Confirm your email address first.")); |
| 325 | } |
| 326 | self.store.set_upkeep(&repo.repo_id, a.enabled, &a.actor.username).await?; |
| 327 | Ok(Outcome::Ok(a.enabled)) |
| 328 | } |
| 329 | |
| 330 | async fn workspace(&self, a: WorkspaceArgs) -> Result<Outcome<Vec<RepoSecurity>>> { |
| 331 | let workspace = a.workspace.to_lowercase(); |
| 332 | if !a.viewer.as_ref().is_some_and(|user| user.is_member(&workspace)) { |
| 333 | return Ok(fail(FailureCode::NotFound, "Workspace not found.")); |
| 334 | } |
| 335 | let mut list = Vec::new(); |
| 336 | for repo in self.store.in_namespace(&workspace).await? { |
| 337 | // Only the repositories whose findings the viewer may see. The |
| 338 | // Write role is never had through being public, so treating |
| 339 | // each as private changes nothing. |
| 340 | let target = access::RepoRef { id: &repo.repo_id, namespace: &workspace, private: true }; |
| 341 | if !access::can(a.viewer.as_ref(), target, SEE_FINDINGS) { |
| 342 | continue; |
| 343 | } |
| 344 | let (counts, secrets, vulnerabilities) = self.store.counts(&repo.repo_id).await?; |
| 345 | list.push(RepoSecurity { |
| 346 | repo_id: repo.repo_id, |
| 347 | name: repo.name, |
| 348 | counts, |
| 349 | secrets, |
| 350 | vulnerabilities, |
| 351 | upkeep: repo.upkeep != 0, |
| 352 | dependencies_scanned_at: repo.deps_scanned_at, |
| 353 | }); |
| 354 | } |
| 355 | Ok(Outcome::Ok(list)) |
| 356 | } |
| 357 | |
| 358 | /// Push protection's question: which of these secrets were allowed? |
| 359 | /// The others are recorded as blocked, so someone can allow them. |
| 360 | async fn push_blocked(&self, a: PushBlockedArgs) -> Result<PushVerdict> { |
| 361 | self.store.register(&a.repo_id, &a.path.namespace, &a.path.name).await?; |
| 362 | let fingerprints: Vec<String> = a.secrets.iter().map(|secret| secret.fingerprint.clone()).collect(); |
| 363 | let known = self.store.known(&a.repo_id, &fingerprints).await?; |
| 364 | let allowed = let_through(&known, &a.secrets); |
| 365 | let fresh: Vec<NewSecret> = a |
| 366 | .secrets |
| 367 | .into_iter() |
| 368 | .filter(|secret| !known.iter().any(|(fingerprint, _, _)| *fingerprint == secret.fingerprint)) |
| 369 | .collect(); |
| 370 | // A likely test value goes through, so it lands: open, not blocked. |
| 371 | let (tests, real): (Vec<NewSecret>, Vec<NewSecret>) = fresh.into_iter().partition(|secret| secret.test_value.is_some()); |
| 372 | self.store |
| 373 | .add_secrets(&a.repo_id, &real, SecretStatus::Blocked, "push", a.pusher.as_deref()) |
| 374 | .await?; |
| 375 | self.store |
| 376 | .add_secrets(&a.repo_id, &tests, SecretStatus::Open, "push", a.pusher.as_deref()) |
| 377 | .await?; |
| 378 | let ids = self |
| 379 | .store |
| 380 | .known(&a.repo_id, &fingerprints) |
| 381 | .await? |
| 382 | .into_iter() |
| 383 | .filter(|(fingerprint, _, _)| !allowed.contains(fingerprint)) |
| 384 | .map(|(fingerprint, id, _)| (fingerprint, id)) |
| 385 | .collect(); |
| 386 | Ok(PushVerdict { allowed, ids }) |
| 387 | } |
| 388 | |
| 389 | /// What happens on the bus that concerns this service. |
| 390 | async fn on_event(&self, event: &Event) -> Result<()> { |
| 391 | match event.kind.as_str() { |
| 392 | "git.push" => { |
| 393 | let Ok(pushed) = serde_json::from_value::<Pushed>(event.data.clone()) else { |
| 394 | return Ok(()); |
| 395 | }; |
| 396 | // Too large to scan before it was stored: its new commits |
| 397 | // are scanned now, on whichever branch. |
| 398 | if pushed.unscanned |
| 399 | && !pushed.after.is_empty() |
| 400 | && let Some(repo) = self.register_by_id(&pushed.repo_id).await? |
| 401 | { |
| 402 | let id = self |
| 403 | .store |
| 404 | .add_push_scan(&repo.repo_id, &pushed.git_ref, &pushed.after, pushed.before.as_deref(), event.actor.as_deref()) |
| 405 | .await?; |
| 406 | self.advance_push_scan(&id, history::PUSH_PAGES_AT_ONCE).await?; |
| 407 | } |
| 408 | // A security update's branch, pushed by its sandbox: time for |
| 409 | // its pull request. |
| 410 | if let Some(branch) = pushed.git_ref.strip_prefix("refs/heads/") |
| 411 | && branch.starts_with(UPDATE_BRANCH_PREFIX) |
| 412 | { |
| 413 | self.update_pushed(&pushed.repo_id, branch).await?; |
| 414 | return Ok(()); |
| 415 | } |
| 416 | if !pushed.default_branch { |
| 417 | return Ok(()); |
| 418 | } |
| 419 | if let Some(repo) = self.register_by_id(&pushed.repo_id).await? { |
| 420 | self.scan_dependencies(&repo).await?; |
| 421 | if repo.history != "done" { |
| 422 | self.advance_history(&repo, 1).await?; |
| 423 | } |
| 424 | } |
| 425 | } |
| 426 | "pull.merged" | "pull.closed" | "checks.completed" => { |
| 427 | if let Ok(happened) = serde_json::from_value::<PullHappened>(event.data.clone()) { |
| 428 | self.update_pull_event( |
| 429 | &event.kind, |
| 430 | &happened.repo_id, |
| 431 | happened.number, |
| 432 | happened.status.as_deref(), |
| 433 | event.actor.as_deref(), |
| 434 | ) |
| 435 | .await?; |
| 436 | } |
| 437 | } |
| 438 | "repo.created" => { |
| 439 | if let Ok(created) = serde_json::from_value::<Created>(event.data.clone()) { |
| 440 | self.register_by_id(&created.repo_id).await?; |
| 441 | } |
| 442 | } |
| 443 | // A repository transferred or renamed: it is recorded under its new path. |
| 444 | "repo.transferred" | "repo.renamed" => { |
| 445 | if let Some(moved) = g1t_kit::transfer::read(event) { |
| 446 | // Where it is now, so moves heard out of order end in |
| 447 | // the same place. |
| 448 | let now: Option<g1t_contracts::repos::RepoPath> = g1t_kit::call( |
| 449 | &self.repos, |
| 450 | "path_by_id", |
| 451 | &g1t_contracts::repos::PathByIdArgs { id: moved.repo_id.clone() }, |
| 452 | ) |
| 453 | .await?; |
| 454 | let current = now.map_or_else( |
| 455 | || moved.destination().to_owned(), |
| 456 | |path| format!("{}/{}", path.namespace, path.name), |
| 457 | ); |
| 458 | if let Some((namespace, name)) = current.split_once('/') { |
| 459 | self.store.moved(&moved.repo_id, namespace, name).await?; |
| 460 | } |
| 461 | } |
| 462 | } |
| 463 | // A repository purged: everything found in it goes. |
| 464 | "repo.purged" => { |
| 465 | if let Some(g1t_kit::lifecycle::Lifecycle::Purged(purged)) = g1t_kit::lifecycle::read(event) { |
| 466 | self.store.purge(&purged.repo_id).await?; |
| 467 | } |
| 468 | } |
| 469 | "workspace.renamed" => { |
| 470 | if let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) { |
| 471 | self.store.rename_namespace(&renamed.stale_slugs(&renamed.to), &renamed.to).await?; |
| 472 | } |
| 473 | } |
| 474 | _ => {} |
| 475 | } |
| 476 | Ok(()) |
| 477 | } |
| 478 | |
| 479 | /// The sweep: continues history scans and scans of pushes that landed |
| 480 | /// unscanned, catches security updates whose sandbox never pushed, and |
| 481 | /// reads dependencies that have not been read for a day. |
| 482 | async fn sweep(&self) -> Result<()> { |
| 483 | for scan in self.store.pending_push_scans(HISTORIES_PER_SWEEP).await? { |
| 484 | if let Err(error) = self.advance_push_scan(&scan.id, PAGES_PER_SWEEP).await { |
| 485 | worker::console_error!("security: push scan {} not continued: {error}", scan.id); |
| 486 | } |
| 487 | } |
| 488 | if let Err(error) = self.stalled_updates().await { |
| 489 | worker::console_error!("security: stalled security updates not handled: {error}"); |
| 490 | } |
| 491 | // Archived and deleted repositories wait; a few more are looked at |
| 492 | // so that they do not hold up the rest. |
| 493 | let mut histories = 0; |
| 494 | for repo in self.store.unfinished_histories(HISTORIES_PER_SWEEP * 4).await? { |
| 495 | if histories == HISTORIES_PER_SWEEP { |
| 496 | break; |
| 497 | } |
| 498 | if !self.active(&repo.repo_id).await? { |
| 499 | continue; |
| 500 | } |
| 501 | histories += 1; |
| 502 | if let Err(error) = self.advance_history(&repo, PAGES_PER_SWEEP).await { |
| 503 | worker::console_error!("security: history of {} not scanned: {error}", repo.repo_id); |
| 504 | } |
| 505 | } |
| 506 | let day_ago = rfc3339(now_ms().saturating_sub(DAY_MS)); |
| 507 | for repo in self.store.stale_dependencies(&day_ago, DEPENDENCIES_PER_SWEEP).await? { |
| 508 | if !self.active(&repo.repo_id).await? { |
| 509 | self.store |
| 510 | .skip_dependencies(&repo.repo_id, "Dependencies are not checked while the repository is archived or deleted.") |
| 511 | .await?; |
| 512 | continue; |
| 513 | } |
| 514 | if let Err(error) = self.scan_dependencies(&repo).await { |
| 515 | worker::console_error!("security: dependencies of {} not read: {error}", repo.repo_id); |
| 516 | } |
| 517 | } |
| 518 | Ok(()) |
| 519 | } |
| 520 | } |
| 521 | |
| 522 | /// The fingerprints a push may carry: those someone dismissed (allowed), |
| 523 | /// and likely test values, which are recorded but never stop a push. |
| 524 | /// `known` is (fingerprint, id, status) of findings already recorded. |
| 525 | fn let_through(known: &[(String, String, String)], secrets: &[NewSecret]) -> Vec<String> { |
| 526 | let mut allowed: Vec<String> = known |
| 527 | .iter() |
| 528 | .filter(|(_, _, status)| status == "allowed") |
| 529 | .map(|(fingerprint, _, _)| fingerprint.clone()) |
| 530 | .chain(secrets.iter().filter(|secret| secret.test_value.is_some()).map(|secret| secret.fingerprint.clone())) |
| 531 | .collect(); |
| 532 | allowed.sort(); |
| 533 | allowed.dedup(); |
| 534 | allowed |
| 535 | } |
| 536 | |
| 537 | #[event(fetch)] |
| 538 | async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> { |
| 539 | let Some(method) = rpc_method(&request) else { |
| 540 | return Response::error("Not found", 404); |
| 541 | }; |
| 542 | let security = Security::new(&env)?; |
| 543 | let body: serde_json::Value = request.json().await?; |
| 544 | match method.as_str() { |
| 545 | "overview" => reply(&security.overview(args(body)?).await?), |
| 546 | "dismiss" => reply(&security.dismiss(args(body)?).await?), |
| 547 | "reopen" => reply(&security.reopen(args(body)?).await?), |
| 548 | "rescan" => reply(&security.rescan(args(body)?).await?), |
| 549 | "set_upkeep" => reply(&security.set_upkeep(args(body)?).await?), |
| 550 | "workspace" => reply(&security.workspace(args(body)?).await?), |
| 551 | "push_blocked" => reply(&security.push_blocked(args(body)?).await?), |
| 552 | _ => Response::error("Unknown method", 404), |
| 553 | } |
| 554 | } |
| 555 | |
| 556 | #[event(queue)] |
| 557 | async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> { |
| 558 | let security = Security::new(&env)?; |
| 559 | for message in batch.messages()? { |
| 560 | let event = message.body(); |
| 561 | if let Err(error) = security.on_event(event).await { |
| 562 | // Scans are idempotent and the sweep catches up, so one failed |
| 563 | // event is logged rather than retried. |
| 564 | worker::console_error!("security: {} {} failed: {error}", event.kind, event.id); |
| 565 | } |
| 566 | } |
| 567 | Ok(()) |
| 568 | } |
| 569 | |
| 570 | #[event(scheduled)] |
| 571 | async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { |
| 572 | match Security::new(&env) { |
| 573 | Ok(security) => { |
| 574 | if let Err(error) = security.sweep().await { |
| 575 | worker::console_error!("security: the sweep failed: {error}"); |
| 576 | } |
| 577 | } |
| 578 | Err(error) => worker::console_error!("security: could not start: {error}"), |
| 579 | } |
| 580 | } |
| 581 | |
| 582 | #[cfg(test)] |
| 583 | mod tests { |
| 584 | use super::*; |
| 585 | |
| 586 | fn secret(fingerprint: &str, test_value: Option<&str>) -> NewSecret { |
| 587 | NewSecret { |
| 588 | fingerprint: fingerprint.into(), |
| 589 | kind: "aws_access_key".into(), |
| 590 | path: "a.env".into(), |
| 591 | line: 1, |
| 592 | commit: "c".into(), |
| 593 | preview: "AKIA…".into(), |
| 594 | test_value: test_value.map(str::to_owned), |
| 595 | } |
| 596 | } |
| 597 | |
| 598 | #[test] |
| 599 | fn dismissed_secrets_and_test_values_go_through() { |
| 600 | let known = vec![ |
| 601 | ("allowed".to_owned(), "sec_1".to_owned(), "allowed".to_owned()), |
| 602 | ("resolved".to_owned(), "sec_2".to_owned(), "resolved".to_owned()), |
| 603 | ("blocked".to_owned(), "sec_3".to_owned(), "blocked".to_owned()), |
| 604 | ]; |
| 605 | let secrets = [secret("allowed", None), secret("resolved", None), secret("example", Some("it says it is an example")), secret("real", None)]; |
| 606 | assert_eq!(let_through(&known, &secrets), ["allowed", "example"]); |
| 607 | } |
| 608 | |
| 609 | #[test] |
| 610 | fn a_push_says_when_it_landed_unscanned() { |
| 611 | let pushed: Pushed = serde_json::from_value(serde_json::json!({ |
| 612 | "repoId": "rep_1", "ref": "refs/heads/import", "after": "abc", "defaultBranch": false, "unscanned": true |
| 613 | })) |
| 614 | .unwrap(); |
| 615 | assert!(pushed.unscanned && pushed.before.is_none() && pushed.git_ref == "refs/heads/import"); |
| 616 | let ordinary: Pushed = serde_json::from_value(serde_json::json!({ "repoId": "rep_1", "ref": "refs/heads/main", "after": "abc", "defaultBranch": true })).unwrap(); |
| 617 | assert!(!ordinary.unscanned); |
| 618 | } |
| 619 | |
| 620 | #[test] |
| 621 | fn owners_are_told_what_landed_and_what_to_do() { |
| 622 | let one = history::landed_secrets_intro("acme", "rocket", "import", 1); |
| 623 | assert!(one.contains("A push to import in acme/rocket") && one.contains("a secret that looks real") && one.contains("Rotate")); |
| 624 | assert!(history::landed_secrets_intro("acme", "rocket", "main", 3).contains("3 secrets that look real")); |
| 625 | } |
| 626 | |
| 627 | #[test] |
| 628 | fn dismissing_a_secret_takes_admin_and_a_dependency_write() { |
| 629 | assert_eq!(Security::capability_for("sec_1"), Capability::ManageIntegrations); |
| 630 | assert_eq!(Security::capability_for("vul_1"), Capability::Push); |
| 631 | } |
| 632 | } |