Skip to content
308 linesCodeBlameRaw

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.

Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails1//! Which events each subscribing service's queue is sent.
2//!
3//! The events service passes each batch on to one queue per subscriber
4//! (`SUBSCRIBER_*` bindings). Each is sent only the types its consumer acts
5//! on, so an event costs a queue write for the services that read it, not
6//! for all of them. A consumer that starts acting on a new type adds it
7//! here in the same change, or it never hears of it; the tests below and
8//! beside each consumer check what they can.
9//!
10//! A pattern is an exact type (`git.push`), a family (`pull.*`, every type
11//! that starts `pull.`), or `*`, every type. A binding not named here is
12//! sent everything, so a new subscriber is never left out.
13
14/// What the shared helpers in `g1t_kit` act on (`rename`, `transfer`,
15/// `lifecycle`, `deleted`, `user_deleted`): rows that follow a workspace
16/// or repository, or go with it. Rare, so every subscriber is sent them.
17const LIFECYCLE: [&str; 8] = [
18 "workspace.renamed",
19 "workspace.deleted",
20 "repo.renamed",
21 "repo.transferred",
22 "repo.deleted",
23 "repo.restored",
24 "repo.purged",
25 "user.deleted",
26];
27
28/// Each subscriber's binding and the types it is sent beside [`LIFECYCLE`].
29pub const ROUTES: &[(&str, &[&str])] = &[
30 // Memory captured from merges and comments, mentions, pull request
31 // bases, runs stopped in archived repositories (work/src/lib.rs queue).
32 (
33 "SUBSCRIBER_WORK",
34 &[
35 "git.push",
36 "branch.renamed",
37 "repo.created",
38 "repo.archived",
39 "repo.default_branch_changed",
40 "pull.merged",
41 "comment.created",
42 ],
43 ),
44 // Pull requests g1t is seeing through, mentions, queued issues and the
45 // merge queue (runner/src/index.ts queue).
46 (
47 "SUBSCRIBER_RUNNER",
48 &[
49 "pull.opened",
50 "pull.ready",
51 "pull.updated",
52 "pull.mergeability",
53 "pull.mergecheck",
54 "pull.merge_requested",
55 "pull.merged",
56 "pull.closed",
57 "checks.completed",
58 "review.completed",
59 "queue.changed",
60 "comment.created",
61 "issue.opened",
62 "issue.updated",
63 "issue.closed",
64 "agent.asked",
65 "repo.archived",
66 ],
67 ),
68 // Push mirrors and write-back to linked trackers (integrations/src).
69 ("SUBSCRIBER_INTEGRATIONS", &["git.push", "issue.closed", "pull.opened"]),
70 // A hook may ask for any type, or for every one.
71 ("SUBSCRIBER_WEBHOOKS", &["*"]),
72 // Every event a workflow can run on (g1t_actions::events::github_events),
73 // and runs stopped in archived repositories.
74 (
75 "SUBSCRIBER_ACTIONS",
76 &[
77 "git.push",
78 "pull.*",
79 "issue.*",
80 "comment.*",
81 "release.*",
82 "deployment.created",
83 "deployment_status.created",
84 "review.completed",
85 "workflow.completed",
86 "repo.archived",
87 ],
88 ),
89 // Previews and production deploys (deployments/src/index.ts queue).
90 (
91 "SUBSCRIBER_DEPLOYMENTS",
92 &[
93 "git.push",
94 "branch.renamed",
95 "repo.default_branch_changed",
96 "pull.opened",
97 "pull.ready",
98 "pull.updated",
99 "pull.reopened",
100 "pull.closed",
101 "pull.merged",
102 "workspace.deleting",
103 "workspace.restored",
104 ],
105 ),
106 // Projects follow their repositories, and wake on activity
107 // (projects/src/index.ts, activity.ts).
108 (
109 "SUBSCRIBER_PROJECTS",
110 &[
111 "git.push",
112 "repo.*",
113 "issue.opened",
114 "issue.closed",
115 "issue.reopened",
116 "pull.opened",
117 "pull.updated",
118 "pull.merged",
119 "pull.closed",
120 "review.completed",
121 "comment.created",
122 "deployment.succeeded",
123 "deployment.failed",
124 "package.published",
125 "package.deleted",
126 "package.version_deleted",
127 "workspace.deleting",
128 "workspace.restored",
129 ],
130 ),
131 // Pull request working copies, workspaces' repositories, contributors
132 // shown as ghost (repos/src/lib.rs handle_event).
133 (
134 "SUBSCRIBER_REPOS",
135 &[
136 "pull.merged",
137 "pull.closed",
138 "pull.reopened",
139 "workspace.deleting",
140 "workspace.restored",
141 "user.deleting",
142 "user.restored",
143 ],
144 ),
145 // Usage rows that follow a renamed workspace or repository.
146 ("SUBSCRIBER_BILLING", &[]),
147 // Secret and dependency scans, and version update pull requests
148 // (security/src/lib.rs on_event).
149 (
150 "SUBSCRIBER_SECURITY",
151 &[
152 "git.push",
153 "pull.opened",
154 "pull.updated",
155 "pull.ready",
156 "pull.merged",
157 "pull.closed",
158 "checks.completed",
159 "comment.created",
160 "repo.created",
161 "repo.visibility_changed",
162 ],
163 ),
164 // What agents are told about a repository (context/src/index.ts).
165 (
166 "SUBSCRIBER_CONTEXT",
167 &[
168 "git.push",
169 "issue.opened",
170 "issue.updated",
171 "issue.closed",
172 "pull.ready",
173 "pull.merged",
174 "memory.changed",
175 ],
176 ),
177 // The site-wide search index (search/src/index.rs).
178 ("SUBSCRIBER_SEARCH", &["git.push", "repo.*", "issue.*", "pull.*", "user.*", "workspace.*"]),
179 // Composer packages read from repositories, and packages that follow
180 // their workspace, team or repository (packages/src/lib.rs on_event).
181 (
182 "SUBSCRIBER_PACKAGES",
183 &[
184 "git.push",
185 "repo.visibility_changed",
186 "workspace.deleting",
187 "workspace.restored",
188 "team.edited",
189 "team.deleted",
190 ],
191 ),
192];
193
194/// Whether `pattern` (see the module's notes) matches the type `kind`.
195fn matches(pattern: &str, kind: &str) -> bool {
196 match pattern.strip_suffix('*') {
197 Some(prefix) => kind.starts_with(prefix),
198 None => pattern == kind,
199 }
200}
201
202/// Whether the subscriber behind `binding` is sent events of type `kind`.
203pub fn routed(binding: &str, kind: &str) -> bool {
204 match ROUTES.iter().find(|(name, _)| *name == binding) {
205 Some((_, patterns)) => {
206 LIFECYCLE.contains(&kind) || patterns.iter().any(|pattern| matches(pattern, kind))
207 }
208 None => true,
209 }
210}
211
212#[cfg(test)]
213mod tests {
214 use super::*;
215
216 /// Types published on the bus that hooks are not offered.
217 const UNOFFERED: [&str; 18] = [
218 "pull.mergecheck",
219 "pull.mergeability",
220 "deployment.review_requested",
221 "memory.changed",
222 "abuse.flagged",
223 "invite.created",
224 "invite.redeemed",
225 "waitlist.requested",
226 "user.updated",
227 "user.deleting",
228 "user.restored",
229 "user.deleted",
230 "user.email_added",
231 "workspace.updated",
232 "workspace.renamed",
233 "workspace.deleting",
234 "workspace.restored",
235 "workspace.deleted",
236 ];
237
238 fn published(kind: &str) -> bool {
239 crate::webhooks::EVENT_TYPES.contains(&kind) || UNOFFERED.contains(&kind)
240 }
241
242 #[test]
243 fn every_exact_route_names_a_type_that_is_published() {
244 for (binding, patterns) in ROUTES {
245 for pattern in patterns.iter().chain(LIFECYCLE.iter()) {
246 if pattern.ends_with('*') {
247 let prefix = pattern.trim_end_matches('*');
248 assert!(
249 prefix.is_empty()
250 || crate::webhooks::EVENT_TYPES.iter().chain(UNOFFERED.iter()).any(|kind| kind.starts_with(prefix)),
251 "{binding}: {pattern} matches nothing"
252 );
253 } else {
254 assert!(published(pattern), "{binding}: {pattern} is not a type anything publishes");
255 }
256 }
257 }
258 }
259
260 #[test]
261 fn every_route_is_a_subscriber_binding_named_once() {
262 let mut names: Vec<&str> = ROUTES.iter().map(|(name, _)| *name).collect();
263 assert!(names.iter().all(|name| name.starts_with("SUBSCRIBER_")));
264 let count = names.len();
265 names.sort_unstable();
266 names.dedup();
267 assert_eq!(names.len(), count);
268 assert_eq!(count, 13);
269 }
270
271 #[test]
272 fn subscribers_hear_what_they_act_on_and_not_the_rest() {
273 assert!(routed("SUBSCRIBER_WEBHOOKS", "secret_scanning_alert.created"));
274 assert!(routed("SUBSCRIBER_WEBHOOKS", "anything.new"));
275 assert!(routed("SUBSCRIBER_ACTIONS", "pull.labeled"));
276 assert!(routed("SUBSCRIBER_ACTIONS", "release.published"));
277 assert!(!routed("SUBSCRIBER_ACTIONS", "check_run.completed"));
278 assert!(routed("SUBSCRIBER_RUNNER", "comment.created"));
279 assert!(!routed("SUBSCRIBER_RUNNER", "status.created"));
280 assert!(routed("SUBSCRIBER_BILLING", "workspace.renamed"));
281 assert!(!routed("SUBSCRIBER_BILLING", "git.push"));
282 assert!(routed("SUBSCRIBER_SEARCH", "user.updated"));
283 assert!(routed("SUBSCRIBER_REPOS", "user.deleting"));
284 assert!(!routed("SUBSCRIBER_REPOS", "git.push"));
285 // Every subscriber follows what moves or removes a repository.
286 for (binding, _) in ROUTES {
287 for kind in LIFECYCLE {
288 assert!(routed(binding, kind), "{binding} misses {kind}");
289 }
290 }
291 }
292
293 #[test]
294 fn a_subscriber_not_in_the_table_hears_everything() {
295 assert!(routed("SUBSCRIBER_NEW", "git.push"));
296 assert!(routed("SUBSCRIBER_NEW", "status.created"));
297 }
298
299 #[test]
300 fn a_family_matches_only_its_own_types() {
301 assert!(matches("pull.*", "pull.opened"));
302 assert!(!matches("pull.*", "pulls.opened"));
303 assert!(!matches("pull.*", "pull"));
304 assert!(matches("*", "x"));
305 assert!(matches("git.push", "git.push"));
306 assert!(!matches("git.push", "git.pushed"));
307 }
308}

This file's history is long; its oldest lines are credited to the oldest commit read.