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

438 lines18,668 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API1//! 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
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look20use g1t_contracts::access::{self, Capability};
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API21use g1t_contracts::events::{Event, WorkspaceRenamed};
22use g1t_contracts::identity::{SlugArgs, Workspace};
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look23use g1t_contracts::repos::{GetArgs, PathByIdArgs, Repo, RepoPath, RepoStatus, StatusByIdArgs};
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API24use g1t_contracts::security::*;
25use g1t_contracts::time::rfc3339;
26use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, User};
27use g1t_kit::{args, now_ms, reply, rpc_method};
28use serde::Deserialize;
29use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
30
31use store::{RepoRow, Store};
32
33/// Repositories whose history is continued per sweep, and pages each.
34const HISTORIES_PER_SWEEP: u32 = 5;
35const PAGES_PER_SWEEP: u32 = 4;
36/// Repositories whose dependencies are read again per sweep.
37const DEPENDENCIES_PER_SWEEP: u32 = 10;
38const DAY_MS: u64 = 24 * 60 * 60 * 1000;
39const MAX_REASON_CHARS: usize = 500;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look40/// What seeing a repository's findings takes: the Write role, as changing
41/// its code does. Dismissing or allowing one takes Admin.
42const SEE_FINDINGS: Capability = Capability::Push;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API43
44pub struct Security {
45 store: Store,
46 identity: Fetcher,
47 repos: Fetcher,
48 work: Fetcher,
49 runner: Fetcher,
50 billing: Fetcher,
51}
52
53fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
54 Outcome::fail(code, message)
55}
56
57/// `git.push`, as far as this service reads it.
58#[derive(Deserialize)]
59#[serde(rename_all = "camelCase")]
60struct Pushed {
61 repo_id: String,
62 #[serde(default)]
63 default_branch: bool,
64}
65
66#[derive(Deserialize)]
67#[serde(rename_all = "camelCase")]
68struct Created {
69 repo_id: String,
70}
71
72impl Security {
73 fn new(env: &Env) -> Result<Self> {
74 Ok(Security {
75 store: Store { db: env.d1("DB")? },
76 identity: env.service("IDENTITY")?,
77 repos: env.service("REPOS")?,
78 work: env.service("WORK")?,
79 runner: env.service("RUNNER")?,
80 billing: env.service("BILLING")?,
81 })
82 }
83
84 /// The workspace itself, acting through no one: who opens upgrade
85 /// issues and asks for an agent on them.
86 async fn workspace_actor(&self, slug: &str) -> Result<Option<User>> {
87 let workspace: Option<Workspace> =
88 g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_owned() }).await?;
89 Ok(workspace.map(|workspace| User {
90 id: workspace.id,
91 username: workspace.slug.clone(),
92 kind: PrincipalKind::Workspace,
93 verified: true,
94 workspaces: vec![Membership::member(workspace.slug)],
95 ..User::default()
96 }))
97 }
98
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look99 /// The repository at `path`, recorded here, if `viewer` may see its
100 /// findings (the Write role) and do `capability`. Findings are for those
101 /// who can change the code: to anyone else the page does not exist,
102 /// public repository or not.
103 async fn member_repo(
104 &self,
105 path: &RepoPath,
106 viewer: &Option<User>,
107 capability: Capability,
108 ) -> Result<Outcome<RepoRow>> {
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API109 let hidden = || fail(FailureCode::NotFound, "Repository not found.");
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look110 if viewer.is_none() {
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API111 return Ok(hidden());
112 }
113 let repo: Outcome<Repo> =
114 g1t_kit::call(&self.repos, "get", &GetArgs { path: path.clone(), viewer: viewer.clone() }).await?;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look115 let repo = match repo {
116 Outcome::Ok(repo) if repo.fork_of.is_none() => repo,
117 _ => return Ok(hidden()),
118 };
119 if !access::can(viewer.as_ref(), &repo, SEE_FINDINGS) {
120 return Ok(hidden());
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API121 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look122 if !access::can(viewer.as_ref(), &repo, capability) {
123 return Ok(fail(
124 FailureCode::Forbidden,
125 access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)),
126 ));
127 }
128 Ok(Outcome::Ok(self.store.register(&repo.id, &repo.namespace, &repo.name).await?))
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API129 }
130
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look131 /// Whether the repository is neither archived nor deleted. When repos
132 /// cannot say, it is taken as active.
133 async fn active(&self, repo_id: &str) -> Result<bool> {
134 let status: Result<RepoStatus> =
135 g1t_kit::call(&self.repos, "status_by_id", &StatusByIdArgs { id: repo_id.to_owned() }).await;
136 Ok(match status {
137 Ok(status) => status.active(),
138 Err(error) => {
139 worker::console_error!("security: status_by_id {repo_id}: {error}");
140 true
141 }
142 })
143 }
144
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API145 /// Records a repository named in an event, by id. Forks are not
146 /// recorded: a pull request's findings belong to its repository.
147 async fn register_by_id(&self, repo_id: &str) -> Result<Option<RepoRow>> {
148 if let Some(row) = self.store.repo(repo_id).await? {
149 return Ok(Some(row));
150 }
151 let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await?;
152 match path {
153 Some(path) => Ok(Some(self.store.register(repo_id, &path.namespace, &path.name).await?)),
154 None => Ok(None),
155 }
156 }
157
158 async fn overview(&self, a: OverviewArgs) -> Result<Outcome<SecurityOverview>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look159 let mut repo = match self.member_repo(&a.repo, &a.viewer, SEE_FINDINGS).await? {
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API160 Outcome::Ok(repo) => repo,
161 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
162 };
163 // The first look at a repository reads its dependencies at once;
164 // its history is scanned in the background.
165 if repo.deps_scanned_at.is_none() {
166 self.scan_dependencies(&repo).await?;
167 repo = self.store.repo(&repo.repo_id).await?.unwrap_or(repo);
168 }
169 let (counts, _, _) = self.store.counts(&repo.repo_id).await?;
170 Ok(Outcome::Ok(SecurityOverview {
171 repo_id: repo.repo_id.clone(),
172 counts,
173 secrets: self.store.secrets(&repo.repo_id).await?,
174 vulnerabilities: self.store.vulnerabilities(&repo.repo_id).await?,
175 scan: repo.scan_state(),
176 upkeep: repo.upkeep != 0,
177 }))
178 }
179
180 async fn decide_secret(&self, a: DecideSecretArgs) -> Result<Outcome<SecretFinding>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look181 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Capability::ManageIntegrations).await? {
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API182 Outcome::Ok(repo) => repo,
183 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
184 };
185 if !a.actor.verified {
186 return Ok(fail(FailureCode::Forbidden, "Confirm your email address first."));
187 }
188 let status = match a.decision.as_str() {
189 "allow" => SecretStatus::Allowed,
190 "resolve" => SecretStatus::Resolved,
191 "reopen" => SecretStatus::Open,
192 _ => return Ok(fail(FailureCode::Invalid, "Say whether to allow, resolve or reopen it.")),
193 };
194 let reason: String = a.reason.trim().chars().take(MAX_REASON_CHARS).collect();
195 if status != SecretStatus::Open && reason.is_empty() {
196 return Ok(fail(
197 FailureCode::Invalid,
198 "Say why: that it is a test fixture, that it was rotated, or that it is meant to be public.",
199 ));
200 }
201 let Some(finding) = self.store.secret(&repo.repo_id, &a.id).await? else {
202 return Ok(fail(FailureCode::NotFound, "No such finding."));
203 };
204 // A secret that never landed has nothing to reopen to but blocked.
205 let status = match (status, finding.status, finding.source.as_str()) {
206 (SecretStatus::Open, _, "push") => SecretStatus::Blocked,
207 (status, _, _) => status,
208 };
209 self.store
210 .decide(&repo.repo_id, &a.id, status, &a.actor.username, Some(&reason))
211 .await?;
212 Ok(match self.store.secret(&repo.repo_id, &a.id).await? {
213 Some(finding) => Outcome::Ok(finding),
214 None => fail(FailureCode::NotFound, "No such finding."),
215 })
216 }
217
218 async fn rescan(&self, a: RescanArgs) -> Result<Outcome<ScanState>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look219 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), SEE_FINDINGS).await? {
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API220 Outcome::Ok(repo) => repo,
221 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
222 };
223 self.store.restart_history(&repo.repo_id).await?;
224 self.scan_dependencies(&repo).await?;
225 if let Some(fresh) = self.store.repo(&repo.repo_id).await? {
226 self.advance_history(&fresh, 1).await?;
227 }
228 let state = self.store.repo(&repo.repo_id).await?.map(|row| row.scan_state()).unwrap_or_default();
229 Ok(Outcome::Ok(state))
230 }
231
232 async fn set_upkeep(&self, a: SetUpkeepArgs) -> Result<Outcome<bool>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look233 // Whether agents keep its dependencies up to date is one of its
234 // settings.
235 let repo = match self.member_repo(&a.repo, &Some(a.actor.clone()), Capability::ManageSettings).await? {
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API236 Outcome::Ok(repo) => repo,
237 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
238 };
239 if !a.actor.verified {
240 return Ok(fail(FailureCode::Forbidden, "Confirm your email address first."));
241 }
242 self.store.set_upkeep(&repo.repo_id, a.enabled, &a.actor.username).await?;
243 Ok(Outcome::Ok(a.enabled))
244 }
245
246 async fn workspace(&self, a: WorkspaceArgs) -> Result<Outcome<Vec<RepoSecurity>>> {
247 let workspace = a.workspace.to_lowercase();
248 if !a.viewer.as_ref().is_some_and(|user| user.is_member(&workspace)) {
249 return Ok(fail(FailureCode::NotFound, "Workspace not found."));
250 }
251 let mut list = Vec::new();
252 for repo in self.store.in_namespace(&workspace).await? {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look253 // Only the repositories whose findings the viewer may see. The
254 // Write role is never had through being public, so treating
255 // each as private changes nothing.
256 let target = access::RepoRef { id: &repo.repo_id, namespace: &workspace, private: true };
257 if !access::can(a.viewer.as_ref(), target, SEE_FINDINGS) {
258 continue;
259 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API260 let (counts, secrets, vulnerabilities) = self.store.counts(&repo.repo_id).await?;
261 list.push(RepoSecurity {
262 repo_id: repo.repo_id,
263 name: repo.name,
264 counts,
265 secrets,
266 vulnerabilities,
267 upkeep: repo.upkeep != 0,
268 dependencies_scanned_at: repo.deps_scanned_at,
269 });
270 }
271 Ok(Outcome::Ok(list))
272 }
273
274 /// Push protection's question: which of these secrets were allowed?
275 /// The others are recorded as blocked, so someone can allow them.
276 async fn push_blocked(&self, a: PushBlockedArgs) -> Result<PushVerdict> {
277 self.store.register(&a.repo_id, &a.path.namespace, &a.path.name).await?;
278 let fingerprints: Vec<String> = a.secrets.iter().map(|secret| secret.fingerprint.clone()).collect();
279 let known = self.store.known(&a.repo_id, &fingerprints).await?;
280 let allowed: Vec<String> = known
281 .iter()
282 .filter(|(_, _, status)| status == "allowed")
283 .map(|(fingerprint, _, _)| fingerprint.clone())
284 .collect();
285 let fresh: Vec<NewSecret> = a
286 .secrets
287 .into_iter()
288 .filter(|secret| !known.iter().any(|(fingerprint, _, _)| *fingerprint == secret.fingerprint))
289 .collect();
290 self.store
291 .add_secrets(&a.repo_id, &fresh, SecretStatus::Blocked, "push", a.pusher.as_deref())
292 .await?;
293 let ids = self
294 .store
295 .known(&a.repo_id, &fingerprints)
296 .await?
297 .into_iter()
298 .filter(|(fingerprint, _, _)| !allowed.contains(fingerprint))
299 .map(|(fingerprint, id, _)| (fingerprint, id))
300 .collect();
301 Ok(PushVerdict { allowed, ids })
302 }
303
304 /// What happens on the bus that concerns this service.
305 async fn on_event(&self, event: &Event) -> Result<()> {
306 match event.kind.as_str() {
307 "git.push" => {
308 let Ok(pushed) = serde_json::from_value::<Pushed>(event.data.clone()) else {
309 return Ok(());
310 };
311 if !pushed.default_branch {
312 return Ok(());
313 }
314 if let Some(repo) = self.register_by_id(&pushed.repo_id).await? {
315 self.scan_dependencies(&repo).await?;
316 if repo.history != "done" {
317 self.advance_history(&repo, 1).await?;
318 }
319 }
320 }
321 "repo.created" => {
322 if let Ok(created) = serde_json::from_value::<Created>(event.data.clone()) {
323 self.register_by_id(&created.repo_id).await?;
324 }
325 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look326 // A repository transferred or renamed: it is recorded under its new path.
327 "repo.transferred" | "repo.renamed" => {
328 if let Some(moved) = g1t_kit::transfer::read(event) {
329 // Where it is now, so moves heard out of order end in
330 // the same place.
331 let now: Option<g1t_contracts::repos::RepoPath> = g1t_kit::call(
332 &self.repos,
333 "path_by_id",
334 &g1t_contracts::repos::PathByIdArgs { id: moved.repo_id.clone() },
335 )
336 .await?;
337 let current = now.map_or_else(
338 || moved.destination().to_owned(),
339 |path| format!("{}/{}", path.namespace, path.name),
340 );
341 if let Some((namespace, name)) = current.split_once('/') {
342 self.store.moved(&moved.repo_id, namespace, name).await?;
343 }
344 }
345 }
346 // A repository purged: everything found in it goes.
347 "repo.purged" => {
348 if let Some(g1t_kit::lifecycle::Lifecycle::Purged(purged)) = g1t_kit::lifecycle::read(event) {
349 self.store.purge(&purged.repo_id).await?;
350 }
351 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API352 "workspace.renamed" => {
353 if let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) {
354 self.store.rename_namespace(&renamed.stale_slugs(&renamed.to), &renamed.to).await?;
355 }
356 }
357 _ => {}
358 }
359 Ok(())
360 }
361
362 /// The sweep: continues history scans, and reads dependencies that
363 /// have not been read for a day.
364 async fn sweep(&self) -> Result<()> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look365 // Archived and deleted repositories wait; a few more are looked at
366 // so that they do not hold up the rest.
367 let mut histories = 0;
368 for repo in self.store.unfinished_histories(HISTORIES_PER_SWEEP * 4).await? {
369 if histories == HISTORIES_PER_SWEEP {
370 break;
371 }
372 if !self.active(&repo.repo_id).await? {
373 continue;
374 }
375 histories += 1;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API376 if let Err(error) = self.advance_history(&repo, PAGES_PER_SWEEP).await {
377 worker::console_error!("security: history of {} not scanned: {error}", repo.repo_id);
378 }
379 }
380 let day_ago = rfc3339(now_ms().saturating_sub(DAY_MS));
381 for repo in self.store.stale_dependencies(&day_ago, DEPENDENCIES_PER_SWEEP).await? {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look382 if !self.active(&repo.repo_id).await? {
383 self.store
384 .skip_dependencies(&repo.repo_id, "Dependencies are not checked while the repository is archived or deleted.")
385 .await?;
386 continue;
387 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API388 if let Err(error) = self.scan_dependencies(&repo).await {
389 worker::console_error!("security: dependencies of {} not read: {error}", repo.repo_id);
390 }
391 }
392 Ok(())
393 }
394}
395
396#[event(fetch)]
397async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
398 let Some(method) = rpc_method(&request) else {
399 return Response::error("Not found", 404);
400 };
401 let security = Security::new(&env)?;
402 let body: serde_json::Value = request.json().await?;
403 match method.as_str() {
404 "overview" => reply(&security.overview(args(body)?).await?),
405 "decide_secret" => reply(&security.decide_secret(args(body)?).await?),
406 "rescan" => reply(&security.rescan(args(body)?).await?),
407 "set_upkeep" => reply(&security.set_upkeep(args(body)?).await?),
408 "workspace" => reply(&security.workspace(args(body)?).await?),
409 "push_blocked" => reply(&security.push_blocked(args(body)?).await?),
410 _ => Response::error("Unknown method", 404),
411 }
412}
413
414#[event(queue)]
415async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
416 let security = Security::new(&env)?;
417 for message in batch.messages()? {
418 let event = message.body();
419 if let Err(error) = security.on_event(event).await {
420 // Scans are idempotent and the sweep catches up, so one failed
421 // event is logged rather than retried.
422 worker::console_error!("security: {} {} failed: {error}", event.kind, event.id);
423 }
424 }
425 Ok(())
426}
427
428#[event(scheduled)]
429async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
430 match Security::new(&env) {
431 Ok(security) => {
432 if let Err(error) = security.sweep().await {
433 worker::console_error!("security: the sweep failed: {error}");
434 }
435 }
436 Err(error) => worker::console_error!("security: could not start: {error}"),
437 }
438}