flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/security/src/lib.rs

632 lines28,088 bytesCodeBlame
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
23mod deps;
24mod history;
25mod security_updates;
26mod store;
27mod updates;
28
29use g1t_contracts::access::{self, Capability};
30use g1t_contracts::events::{Event, WorkspaceRenamed};
31use g1t_contracts::security::UPDATE_BRANCH_PREFIX;
32use g1t_contracts::repos::{GetArgs, PathByIdArgs, Repo, RepoPath, RepoStatus, StatusByIdArgs};
33use g1t_contracts::security::*;
34use g1t_contracts::time::rfc3339;
35use g1t_contracts::{FailureCode, Outcome, User};
36use g1t_kit::{args, now_ms, reply, rpc_method};
37use serde::Deserialize;
38use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
39
40use store::{Activity, RepoRow, Store};
41
42/// Repositories whose history is continued per sweep, and pages each.
43const HISTORIES_PER_SWEEP: u32 = 5;
44const PAGES_PER_SWEEP: u32 = 4;
45/// Repositories whose dependencies are read again per sweep.
46const DEPENDENCIES_PER_SWEEP: u32 = 10;
47const DAY_MS: u64 = 24 * 60 * 60 * 1000;
48const 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.
51const SEE_FINDINGS: Capability = Capability::Push;
52
53pub struct Security {
54 store: Store,
55 identity: Fetcher,
56 repos: Fetcher,
57 work: Fetcher,
58 runner: Fetcher,
59 billing: Fetcher,
60}
61
62fn 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")]
69struct 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")]
88struct PullHappened {
89 repo_id: String,
90 number: u32,
91 #[serde(default)]
92 status: Option<String>,
93}
94
95/// The longest comment a dismissal keeps.
96const MAX_COMMENT_CHARS: usize = MAX_REASON_CHARS;
97/// Activity rows the Security page reads, newest first.
98const ACTIVITY_SHOWN: u32 = 500;
99
100#[derive(Deserialize)]
101#[serde(rename_all = "camelCase")]
102struct Created {
103 repo_id: String,
104}
105
106impl 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.
525fn 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)]
538async 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)]
557async 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)]
571async 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)]
583mod 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}