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

356 lines14,903 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//! - Dependencies. On every push to a default branch, and daily, the
8//! lockfiles are read and every package checked against OSV. Each
9//! vulnerable package with a fix gets one upgrade issue, which a g1t
10//! agent takes through the usual pull request, checks, review and merge
11//! queue.
12//!
13//! Other services reach it over `POST /rpc/<method>`; see
14//! `g1t_contracts::security`.
15
16mod deps;
17mod history;
18mod store;
19
20use g1t_contracts::events::{Event, WorkspaceRenamed};
21use g1t_contracts::identity::{SlugArgs, Workspace};
22use g1t_contracts::repos::{GetArgs, PathByIdArgs, Repo, RepoPath};
23use g1t_contracts::security::*;
24use g1t_contracts::time::rfc3339;
25use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, User};
26use g1t_kit::{args, now_ms, reply, rpc_method};
27use serde::Deserialize;
28use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
29
30use store::{RepoRow, Store};
31
32/// Repositories whose history is continued per sweep, and pages each.
33const HISTORIES_PER_SWEEP: u32 = 5;
34const PAGES_PER_SWEEP: u32 = 4;
35/// Repositories whose dependencies are read again per sweep.
36const DEPENDENCIES_PER_SWEEP: u32 = 10;
37const DAY_MS: u64 = 24 * 60 * 60 * 1000;
38const MAX_REASON_CHARS: usize = 500;
39
40pub struct Security {
41 store: Store,
42 identity: Fetcher,
43 repos: Fetcher,
44 work: Fetcher,
45 runner: Fetcher,
46 billing: Fetcher,
47}
48
49fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
50 Outcome::fail(code, message)
51}
52
53/// `git.push`, as far as this service reads it.
54#[derive(Deserialize)]
55#[serde(rename_all = "camelCase")]
56struct Pushed {
57 repo_id: String,
58 #[serde(default)]
59 default_branch: bool,
60}
61
62#[derive(Deserialize)]
63#[serde(rename_all = "camelCase")]
64struct Created {
65 repo_id: String,
66}
67
68impl Security {
69 fn new(env: &Env) -> Result<Self> {
70 Ok(Security {
71 store: Store { db: env.d1("DB")? },
72 identity: env.service("IDENTITY")?,
73 repos: env.service("REPOS")?,
74 work: env.service("WORK")?,
75 runner: env.service("RUNNER")?,
76 billing: env.service("BILLING")?,
77 })
78 }
79
80 /// The workspace itself, acting through no one: who opens upgrade
81 /// issues and asks for an agent on them.
82 async fn workspace_actor(&self, slug: &str) -> Result<Option<User>> {
83 let workspace: Option<Workspace> =
84 g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_owned() }).await?;
85 Ok(workspace.map(|workspace| User {
86 id: workspace.id,
87 username: workspace.slug.clone(),
88 kind: PrincipalKind::Workspace,
89 verified: true,
90 workspaces: vec![Membership::member(workspace.slug)],
91 ..User::default()
92 }))
93 }
94
95 /// The repository at `path`, recorded here, if `viewer` is a member of
96 /// its workspace. Findings are the workspace's own business: to anyone
97 /// else the page does not exist, public repository or not.
98 async fn member_repo(&self, path: &RepoPath, viewer: &Option<User>) -> Result<Outcome<RepoRow>> {
99 let namespace = path.namespace.to_lowercase();
100 let hidden = || fail(FailureCode::NotFound, "Repository not found.");
101 if !viewer.as_ref().is_some_and(|user| user.is_member(&namespace)) {
102 return Ok(hidden());
103 }
104 let repo: Outcome<Repo> =
105 g1t_kit::call(&self.repos, "get", &GetArgs { path: path.clone(), viewer: viewer.clone() }).await?;
106 match repo {
107 Outcome::Ok(repo) if repo.fork_of.is_none() => {
108 Ok(Outcome::Ok(self.store.register(&repo.id, &repo.namespace, &repo.name).await?))
109 }
110 _ => Ok(hidden()),
111 }
112 }
113
114 /// Records a repository named in an event, by id. Forks are not
115 /// recorded: a pull request's findings belong to its repository.
116 async fn register_by_id(&self, repo_id: &str) -> Result<Option<RepoRow>> {
117 if let Some(row) = self.store.repo(repo_id).await? {
118 return Ok(Some(row));
119 }
120 let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await?;
121 match path {
122 Some(path) => Ok(Some(self.store.register(repo_id, &path.namespace, &path.name).await?)),
123 None => Ok(None),
124 }
125 }
126
127 async fn overview(&self, a: OverviewArgs) -> Result<Outcome<SecurityOverview>> {
128 let mut repo = match self.member_repo(&a.repo, &a.viewer).await? {
129 Outcome::Ok(repo) => repo,
130 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
131 };
132 // The first look at a repository reads its dependencies at once;
133 // its history is scanned in the background.
134 if repo.deps_scanned_at.is_none() {
135 self.scan_dependencies(&repo).await?;
136 repo = self.store.repo(&repo.repo_id).await?.unwrap_or(repo);
137 }
138 let (counts, _, _) = self.store.counts(&repo.repo_id).await?;
139 Ok(Outcome::Ok(SecurityOverview {
140 repo_id: repo.repo_id.clone(),
141 counts,
142 secrets: self.store.secrets(&repo.repo_id).await?,
143 vulnerabilities: self.store.vulnerabilities(&repo.repo_id).await?,
144 scan: repo.scan_state(),
145 upkeep: repo.upkeep != 0,
146 }))
147 }
148
149 async fn decide_secret(&self, a: DecideSecretArgs) -> Result<Outcome<SecretFinding>> {
150 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone())).await? {
151 Outcome::Ok(repo) => repo,
152 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
153 };
154 if !a.actor.verified {
155 return Ok(fail(FailureCode::Forbidden, "Confirm your email address first."));
156 }
157 let status = match a.decision.as_str() {
158 "allow" => SecretStatus::Allowed,
159 "resolve" => SecretStatus::Resolved,
160 "reopen" => SecretStatus::Open,
161 _ => return Ok(fail(FailureCode::Invalid, "Say whether to allow, resolve or reopen it.")),
162 };
163 let reason: String = a.reason.trim().chars().take(MAX_REASON_CHARS).collect();
164 if status != SecretStatus::Open && reason.is_empty() {
165 return Ok(fail(
166 FailureCode::Invalid,
167 "Say why: that it is a test fixture, that it was rotated, or that it is meant to be public.",
168 ));
169 }
170 let Some(finding) = self.store.secret(&repo.repo_id, &a.id).await? else {
171 return Ok(fail(FailureCode::NotFound, "No such finding."));
172 };
173 // A secret that never landed has nothing to reopen to but blocked.
174 let status = match (status, finding.status, finding.source.as_str()) {
175 (SecretStatus::Open, _, "push") => SecretStatus::Blocked,
176 (status, _, _) => status,
177 };
178 self.store
179 .decide(&repo.repo_id, &a.id, status, &a.actor.username, Some(&reason))
180 .await?;
181 Ok(match self.store.secret(&repo.repo_id, &a.id).await? {
182 Some(finding) => Outcome::Ok(finding),
183 None => fail(FailureCode::NotFound, "No such finding."),
184 })
185 }
186
187 async fn rescan(&self, a: RescanArgs) -> Result<Outcome<ScanState>> {
188 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone())).await? {
189 Outcome::Ok(repo) => repo,
190 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
191 };
192 self.store.restart_history(&repo.repo_id).await?;
193 self.scan_dependencies(&repo).await?;
194 if let Some(fresh) = self.store.repo(&repo.repo_id).await? {
195 self.advance_history(&fresh, 1).await?;
196 }
197 let state = self.store.repo(&repo.repo_id).await?.map(|row| row.scan_state()).unwrap_or_default();
198 Ok(Outcome::Ok(state))
199 }
200
201 async fn set_upkeep(&self, a: SetUpkeepArgs) -> Result<Outcome<bool>> {
202 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone())).await? {
203 Outcome::Ok(repo) => repo,
204 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
205 };
206 if !a.actor.verified {
207 return Ok(fail(FailureCode::Forbidden, "Confirm your email address first."));
208 }
209 self.store.set_upkeep(&repo.repo_id, a.enabled, &a.actor.username).await?;
210 Ok(Outcome::Ok(a.enabled))
211 }
212
213 async fn workspace(&self, a: WorkspaceArgs) -> Result<Outcome<Vec<RepoSecurity>>> {
214 let workspace = a.workspace.to_lowercase();
215 if !a.viewer.as_ref().is_some_and(|user| user.is_member(&workspace)) {
216 return Ok(fail(FailureCode::NotFound, "Workspace not found."));
217 }
218 let mut list = Vec::new();
219 for repo in self.store.in_namespace(&workspace).await? {
220 let (counts, secrets, vulnerabilities) = self.store.counts(&repo.repo_id).await?;
221 list.push(RepoSecurity {
222 repo_id: repo.repo_id,
223 name: repo.name,
224 counts,
225 secrets,
226 vulnerabilities,
227 upkeep: repo.upkeep != 0,
228 dependencies_scanned_at: repo.deps_scanned_at,
229 });
230 }
231 Ok(Outcome::Ok(list))
232 }
233
234 /// Push protection's question: which of these secrets were allowed?
235 /// The others are recorded as blocked, so someone can allow them.
236 async fn push_blocked(&self, a: PushBlockedArgs) -> Result<PushVerdict> {
237 self.store.register(&a.repo_id, &a.path.namespace, &a.path.name).await?;
238 let fingerprints: Vec<String> = a.secrets.iter().map(|secret| secret.fingerprint.clone()).collect();
239 let known = self.store.known(&a.repo_id, &fingerprints).await?;
240 let allowed: Vec<String> = known
241 .iter()
242 .filter(|(_, _, status)| status == "allowed")
243 .map(|(fingerprint, _, _)| fingerprint.clone())
244 .collect();
245 let fresh: Vec<NewSecret> = a
246 .secrets
247 .into_iter()
248 .filter(|secret| !known.iter().any(|(fingerprint, _, _)| *fingerprint == secret.fingerprint))
249 .collect();
250 self.store
251 .add_secrets(&a.repo_id, &fresh, SecretStatus::Blocked, "push", a.pusher.as_deref())
252 .await?;
253 let ids = self
254 .store
255 .known(&a.repo_id, &fingerprints)
256 .await?
257 .into_iter()
258 .filter(|(fingerprint, _, _)| !allowed.contains(fingerprint))
259 .map(|(fingerprint, id, _)| (fingerprint, id))
260 .collect();
261 Ok(PushVerdict { allowed, ids })
262 }
263
264 /// What happens on the bus that concerns this service.
265 async fn on_event(&self, event: &Event) -> Result<()> {
266 match event.kind.as_str() {
267 "git.push" => {
268 let Ok(pushed) = serde_json::from_value::<Pushed>(event.data.clone()) else {
269 return Ok(());
270 };
271 if !pushed.default_branch {
272 return Ok(());
273 }
274 if let Some(repo) = self.register_by_id(&pushed.repo_id).await? {
275 self.scan_dependencies(&repo).await?;
276 if repo.history != "done" {
277 self.advance_history(&repo, 1).await?;
278 }
279 }
280 }
281 "repo.created" => {
282 if let Ok(created) = serde_json::from_value::<Created>(event.data.clone()) {
283 self.register_by_id(&created.repo_id).await?;
284 }
285 }
286 "workspace.renamed" => {
287 if let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) {
288 self.store.rename_namespace(&renamed.stale_slugs(&renamed.to), &renamed.to).await?;
289 }
290 }
291 _ => {}
292 }
293 Ok(())
294 }
295
296 /// The sweep: continues history scans, and reads dependencies that
297 /// have not been read for a day.
298 async fn sweep(&self) -> Result<()> {
299 for repo in self.store.unfinished_histories(HISTORIES_PER_SWEEP).await? {
300 if let Err(error) = self.advance_history(&repo, PAGES_PER_SWEEP).await {
301 worker::console_error!("security: history of {} not scanned: {error}", repo.repo_id);
302 }
303 }
304 let day_ago = rfc3339(now_ms().saturating_sub(DAY_MS));
305 for repo in self.store.stale_dependencies(&day_ago, DEPENDENCIES_PER_SWEEP).await? {
306 if let Err(error) = self.scan_dependencies(&repo).await {
307 worker::console_error!("security: dependencies of {} not read: {error}", repo.repo_id);
308 }
309 }
310 Ok(())
311 }
312}
313
314#[event(fetch)]
315async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
316 let Some(method) = rpc_method(&request) else {
317 return Response::error("Not found", 404);
318 };
319 let security = Security::new(&env)?;
320 let body: serde_json::Value = request.json().await?;
321 match method.as_str() {
322 "overview" => reply(&security.overview(args(body)?).await?),
323 "decide_secret" => reply(&security.decide_secret(args(body)?).await?),
324 "rescan" => reply(&security.rescan(args(body)?).await?),
325 "set_upkeep" => reply(&security.set_upkeep(args(body)?).await?),
326 "workspace" => reply(&security.workspace(args(body)?).await?),
327 "push_blocked" => reply(&security.push_blocked(args(body)?).await?),
328 _ => Response::error("Unknown method", 404),
329 }
330}
331
332#[event(queue)]
333async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
334 let security = Security::new(&env)?;
335 for message in batch.messages()? {
336 let event = message.body();
337 if let Err(error) = security.on_event(event).await {
338 // Scans are idempotent and the sweep catches up, so one failed
339 // event is logged rather than retried.
340 worker::console_error!("security: {} {} failed: {error}", event.kind, event.id);
341 }
342 }
343 Ok(())
344}
345
346#[event(scheduled)]
347async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
348 match Security::new(&env) {
349 Ok(security) => {
350 if let Err(error) = security.sweep().await {
351 worker::console_error!("security: the sweep failed: {error}");
352 }
353 }
354 Err(error) => worker::console_error!("security: could not start: {error}"),
355 }
356}