g1t/services/webhooks/src/deliver.rs

286 lines12,427 bytesCodeBlame
1//! What a delivery is made of, apart from sending it: the payload, the
2//! signature, when to try again, and which addresses may be sent to.
3
4use g1t_contracts::events::Event;
5use g1t_contracts::webhooks::EVENT_TYPES;
6use serde_json::{Value, json};
7
8/// How long to wait after each failed attempt before the next, so a
9/// delivery gets six attempts over about seven and a half hours.
10pub const RETRY_WAITS_SECONDS: [u64; 5] = [60, 5 * 60, 30 * 60, 2 * 60 * 60, 5 * 60 * 60];
11
12/// When to try again after `attempts` attempts, in seconds from now, or
13/// `None` when it has had them all.
14pub fn retry_after(attempts: u32) -> Option<u64> {
15 RETRY_WAITS_SECONDS.get(attempts.checked_sub(1)? as usize).copied()
16}
17
18/// Whether a webhook that wants `wanted` is sent an event of type `kind`.
19pub fn wants(wanted: &[String], kind: &str) -> bool {
20 wanted.iter().any(|event| event == "*" || event == kind)
21}
22
23/// Event types as given, checked and in catalogue order. `["*"]` when none
24/// or all are given. `Err` names the first that is not one.
25pub fn tidy_events(given: &[String]) -> Result<Vec<String>, String> {
26 if given.is_empty() || given.iter().any(|event| event == "*") {
27 return Ok(vec!["*".to_owned()]);
28 }
29 if let Some(unknown) = given.iter().find(|event| !EVENT_TYPES.contains(&event.as_str())) {
30 return Err(format!("There is no event called {unknown}."));
31 }
32 Ok(EVENT_TYPES
33 .iter()
34 .filter(|kind| given.iter().any(|event| event == *kind))
35 .map(|kind| (*kind).to_owned())
36 .collect())
37}
38
39/// What is sent for an event: the event as the bus has it, with the
40/// repository, the workspace and whoever caused it named. Its keys are in
41/// `snake_case`, as everything g1t sends out is (see `g1t_kit::wire`).
42pub fn payload(event: &Event, workspace: &str, repo: Option<(&str, &str)>, actor_name: Option<&str>) -> Value {
43 g1t_kit::wire::snake_case(json!({
44 "id": event.id,
45 "type": event.kind,
46 "time": event.time,
47 "workspace": workspace,
48 "repository": repo.map(|(id, full_name)| json!({ "id": id, "full_name": full_name })),
49 "actor": event.actor.as_ref().map(|id| json!({ "id": id, "username": actor_name })),
50 "data": event.data,
51 }))
52}
53
54/// The repository a repository event names, as `(namespace, name)`, where
55/// it is now: for when repos no longer shows it (it was deleted or purged).
56pub fn named_in(event: &Event) -> Option<(String, String)> {
57 let text = |key: &str| event.data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
58 match event.kind.as_str() {
59 "repo.renamed" => Some((text("namespace")?, text("to")?)),
60 "repo.transferred" => Some((text("to")?, text("name")?)),
61 kind if kind.starts_with("repo.") => Some((text("namespace")?, text("name")?)),
62 _ => None,
63 }
64}
65
66/// The workspace an event belongs to when it is about no repository: a
67/// package of the workspace's own, unlinked from any repository. Such an
68/// event goes to the workspace's webhooks only. Events about a repository,
69/// and every other kind, are `None`: they are routed by their repository.
70pub fn workspace_scoped(event: &Event) -> Option<String> {
71 if event.repo_id.as_deref().is_some_and(|id| !id.is_empty()) || !event.kind.starts_with("package.") {
72 return None;
73 }
74 event.data["workspace"]
75 .as_str()
76 .map(|slug| slug.trim().to_lowercase())
77 .filter(|slug| !slug.is_empty())
78}
79
80/// What is sent to check a webhook works.
81pub fn ping(hook_id: &str, url: &str, events: &[String], time: &str) -> Value {
82 json!({
83 "type": "ping",
84 "time": time,
85 "hook": { "id": hook_id, "url": url, "events": events },
86 "message": "g1t will send this webhook's events here.",
87 })
88}
89
90/// The signature header's value: the body's HMAC-SHA256 under the secret.
91pub fn signature(secret: &str, body: &str) -> String {
92 format!("sha256={}", g1t_secrets::hmac_sha256_hex(secret, body))
93}
94
95/// Whether g1t may send to `url`: HTTPS, to a public host. `Err` says why
96/// not.
97pub fn check_url(url: &str) -> Result<(), String> {
98 let Some(rest) = url.strip_prefix("https://") else {
99 return Err("Webhooks are sent over HTTPS: the address must start with https://.".to_owned());
100 };
101 if url.len() > 2000 {
102 return Err("That address is too long.".to_owned());
103 }
104 let authority = rest.split(['/', '?', '#']).next().unwrap_or_default();
105 let host = authority.rsplit('@').next().unwrap_or_default();
106 let host = host.strip_prefix('[').map_or_else(
107 || host.split(':').next().unwrap_or_default(),
108 |v6| v6.split(']').next().unwrap_or_default(),
109 );
110 let host = host.to_ascii_lowercase();
111 if host.is_empty() || !host.contains(['.', ':']) {
112 return Err("The address needs a host g1t can reach on the internet.".to_owned());
113 }
114 let private_name = host == "localhost"
115 || [".localhost", ".local", ".internal", ".lan", ".home.arpa"]
116 .iter()
117 .any(|suffix| host.ends_with(suffix));
118 let private_ip = host.parse::<std::net::IpAddr>().is_ok_and(|ip| match ip {
119 std::net::IpAddr::V4(v4) => {
120 v4.is_private() || v4.is_loopback() || v4.is_link_local() || v4.is_unspecified() || v4.is_broadcast() || v4.octets()[0] == 100 && (64..128).contains(&v4.octets()[1])
121 }
122 std::net::IpAddr::V6(v6) => v6.is_loopback() || v6.is_unspecified() || (v6.segments()[0] & 0xfe00) == 0xfc00 || (v6.segments()[0] & 0xffc0) == 0xfe80,
123 });
124 if private_name || private_ip {
125 return Err("Webhooks go to public addresses, not private or local ones.".to_owned());
126 }
127 Ok(())
128}
129
130#[cfg(test)]
131mod tests {
132 use super::*;
133
134 #[test]
135 fn retries_wait_longer_each_time_then_stop() {
136 assert_eq!(retry_after(1), Some(60));
137 assert_eq!(retry_after(2), Some(300));
138 assert_eq!(retry_after(5), Some(18_000));
139 assert_eq!(retry_after(6), None);
140 assert_eq!(retry_after(0), None);
141 }
142
143 #[test]
144 fn events_are_checked_and_ordered() {
145 let given = vec!["pull.merged".to_owned(), "git.push".to_owned()];
146 assert_eq!(tidy_events(&given).unwrap(), vec!["git.push", "pull.merged"]);
147 assert_eq!(tidy_events(&[]).unwrap(), vec!["*"]);
148 assert!(tidy_events(&["pull.exploded".to_owned()]).is_err());
149 assert!(wants(&["*".to_owned()], "issue.opened"));
150 assert!(!wants(&["git.push".to_owned()], "issue.opened"));
151 }
152
153 #[test]
154 fn only_public_https_addresses_are_allowed() {
155 assert!(check_url("https://hooks.example.com/g1t").is_ok());
156 assert!(check_url("https://example.com:8443/hook?x=1").is_ok());
157 for url in [
158 "http://example.com/hook",
159 "https://localhost/hook",
160 "https://127.0.0.1/hook",
161 "https://10.0.0.5/hook",
162 "https://192.168.1.2:8080/hook",
163 "https://169.254.169.254/latest",
164 "https://100.64.0.1/hook",
165 "https://[::1]/hook",
166 "https://[fd00::1]/hook",
167 "https://printer.local/hook",
168 "https://intranet/hook",
169 "https://user@10.0.0.1/hook",
170 ] {
171 assert!(check_url(url).is_err(), "{url}");
172 }
173 }
174
175 #[test]
176 fn payloads_are_snake_case() {
177 let event = Event {
178 id: "evt_1".to_owned(),
179 kind: "pull.merged".to_owned(),
180 source: "work".to_owned(),
181 time: "2026-10-04T16:00:00Z".to_owned(),
182 repo_id: Some("rep_1".to_owned()),
183 actor: Some("usr_1".to_owned()),
184 data: json!({ "pullId": "pul_1", "repoId": "rep_1", "supersededBy": null, "inputs": { "dryRun": true } }),
185 };
186 let sent = payload(&event, "acme", Some(("rep_1", "acme/rocket")), Some("syntaqx"));
187 assert_eq!(sent["repository"]["full_name"], "acme/rocket");
188 assert_eq!(sent["data"]["pull_id"], "pul_1");
189 assert!(sent["data"].get("superseded_by").is_some());
190 assert_eq!(sent["data"]["inputs"]["dryRun"], true);
191 assert!(g1t_kit::wire::camel_case_keys(&sent).is_empty());
192 }
193
194 #[test]
195 fn a_pull_request_g1t_made_names_g1t_and_who_asked() {
196 // As the work service publishes it: g1t the author, the person who
197 // asked beside it, and the actor whoever caused the event.
198 let event = Event {
199 id: "evt_1".to_owned(),
200 kind: "pull.opened".to_owned(),
201 source: "work".to_owned(),
202 time: "2026-10-06T10:00:00Z".to_owned(),
203 repo_id: Some("rep_1".to_owned()),
204 actor: Some("usr_1".to_owned()),
205 data: json!({
206 "pullId": "pr_1", "repoId": "rep_1", "number": 14, "agent": "g1t",
207 "author": { "id": "usr_g1t_agent", "username": "g1t" },
208 "requestedBy": { "id": "usr_1", "username": "syntaqx" }
209 }),
210 };
211 let sent = payload(&event, "acme", Some(("rep_1", "acme/rocket")), Some("syntaqx"));
212 assert_eq!(sent["data"]["author"]["username"], "g1t");
213 assert_eq!(sent["data"]["requested_by"]["username"], "syntaqx");
214 assert_eq!(sent["actor"]["username"], "syntaqx");
215 }
216
217 fn repo_event(kind: &str, data: Value) -> Event {
218 Event {
219 id: "evt_1".to_owned(),
220 kind: kind.to_owned(),
221 source: "repos".to_owned(),
222 time: "2026-10-05T16:00:00Z".to_owned(),
223 repo_id: Some("rep_1".to_owned()),
224 actor: None,
225 data,
226 }
227 }
228
229 #[test]
230 fn a_workspaces_own_package_events_go_to_its_webhooks() {
231 let mut event = repo_event("package.published", json!({ "packageId": "pkg_1", "workspace": "Acme", "name": "tools", "repoId": null }));
232 event.repo_id = None;
233 assert_eq!(workspace_scoped(&event).as_deref(), Some("acme"));
234 let sent = payload(&event, "acme", None, Some("ana"));
235 assert_eq!(sent["repository"], Value::Null);
236 assert_eq!(sent["workspace"], "acme");
237 assert_eq!(sent["data"]["package_id"], "pkg_1", "snake_case as everything sent");
238 // A linked package's events go by its repository.
239 let linked = repo_event("package.published", json!({ "workspace": "acme", "repoId": "rep_1" }));
240 assert_eq!(workspace_scoped(&linked), None);
241 // Other events without a repository are not workspace events.
242 let mut other = repo_event("issue.opened", json!({ "workspace": "acme" }));
243 other.repo_id = None;
244 assert_eq!(workspace_scoped(&other), None);
245 let mut nameless = repo_event("package.deleted", json!({ "workspace": "" }));
246 nameless.repo_id = None;
247 assert_eq!(workspace_scoped(&nameless), None);
248 }
249
250 #[test]
251 fn repository_events_name_where_it_is_now() {
252 let named = |kind: &str, data: Value| named_in(&repo_event(kind, data));
253 let at = |namespace: &str, name: &str| Some((namespace.to_owned(), name.to_owned()));
254 assert_eq!(named("repo.renamed", json!({ "namespace": "acme", "from": "old", "to": "new" })), at("acme", "new"));
255 assert_eq!(named("repo.transferred", json!({ "name": "web", "from": "a", "to": "b" })), at("b", "web"));
256 assert_eq!(named("repo.purged", json!({ "repoId": "rep_1", "namespace": "acme", "name": "web" })), at("acme", "web"));
257 assert_eq!(named("repo.deleted", json!({ "repoId": "rep_1", "namespace": "", "name": "web" })), None);
258 assert_eq!(named("issue.opened", json!({ "namespace": "acme", "name": "web" })), None);
259 }
260
261 #[test]
262 fn repository_lifecycle_events_can_be_chosen() {
263 for kind in [
264 "repo.updated",
265 "repo.visibility_changed",
266 "repo.renamed",
267 "repo.transferred",
268 "repo.archived",
269 "repo.unarchived",
270 "repo.deleted",
271 "repo.restored",
272 "repo.purged",
273 "repo.default_branch_changed",
274 "branch.renamed",
275 ] {
276 assert_eq!(tidy_events(&[kind.to_owned()]).unwrap(), vec![kind], "{kind}");
277 }
278 }
279
280 #[test]
281 fn the_signature_is_the_bodys_hmac() {
282 let signature = signature("shh", "{\"a\":1}");
283 assert!(signature.starts_with("sha256="));
284 assert!(g1t_secrets::signed("shh", "{\"a\":1}", &signature));
285 }
286}