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.

Each subscriber is sent only the events it acts on, in batches the queues take; failed messages go to a dead-letter queue, repos and packages retry one event at a time, and a push of many tags or branches announces few (events 0007)1//! 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}