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/webhooks/src/lib.rs

805 lines32,026 bytesCodeBlame
1//! The webhooks service: events, delivered to the addresses a repository or
2//! a workspace registers. See `g1t_contracts::webhooks` for the methods and
3//! their arguments.
4//!
5//! An event from the bus becomes a delivery for each active webhook that
6//! wants it, and is sent at once. One that is not answered with a 2xx is
7//! tried again by the minute's sweep, waiting longer each time, until it
8//! has had every attempt. Every delivery is kept for a fortnight, with what
9//! was sent and what came back.
10
11mod deliver;
12mod rename;
13
14use std::time::Duration;
15
16use futures_util::future::{Either, select};
17use g1t_contracts::access::{self, Capability};
18use g1t_contracts::events::Event;
19use g1t_contracts::identity::{AGENT_ID, AGENT_NAME, UsernamesArgs};
20use g1t_contracts::repos::{GetArgs, GetByIdArgs, Repo};
21use g1t_contracts::time::rfc3339;
22use g1t_contracts::webhooks::*;
23use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, Role, User, Viewer, new_id};
24use g1t_kit::{args, now_ms, reply, rpc_method};
25use g1t_secrets::Sealer;
26use serde::Deserialize;
27use serde_json::Value;
28use worker::wasm_bindgen::JsValue;
29use worker::{
30 Context, D1Database, Delay, Env, Fetch, Fetcher, Headers, MessageBatch, MessageExt, Method, Request, RequestInit, Response,
31 Result, ScheduleContext, ScheduledEvent, event,
32};
33
34/// How long a receiver has to answer.
35const TIMEOUT: Duration = Duration::from_secs(10);
36/// How much of an answer is kept.
37const RESPONSE_KEPT: usize = 2_000;
38const DELIVERIES_SHOWN: u32 = 50;
39/// How long deliveries are kept.
40const KEPT_DAYS: u64 = 14;
41/// How many due retries one sweep makes.
42const SWEEP: u32 = 50;
43
44#[derive(Deserialize)]
45struct HookRow {
46 id: String,
47 scope: String,
48 workspace: String,
49 repo: Option<String>,
50 url: String,
51 events: String,
52 active: u32,
53 secret: String,
54 secret_hint: String,
55 created_by: String,
56 created_at: String,
57 last_status: Option<String>,
58 last_delivered_at: Option<String>,
59}
60
61impl HookRow {
62 fn events(&self) -> Vec<String> {
63 serde_json::from_str(&self.events).unwrap_or_else(|_| vec!["*".to_owned()])
64 }
65
66 fn to_hook(&self) -> Hook {
67 Hook {
68 id: self.id.clone(),
69 scope: if self.scope == "repo" { HookScope::Repo } else { HookScope::Workspace },
70 workspace: self.workspace.clone(),
71 repo: self.repo.clone(),
72 url: self.url.clone(),
73 events: self.events(),
74 active: self.active != 0,
75 secret_hint: self.secret_hint.clone(),
76 created_by: self.created_by.clone(),
77 created_at: self.created_at.clone(),
78 last_status: self.last_status.clone(),
79 last_delivered_at: self.last_delivered_at.clone(),
80 }
81 }
82}
83
84#[derive(Deserialize)]
85struct DeliveryRow {
86 id: String,
87 hook_id: String,
88 event_id: String,
89 event: String,
90 payload: String,
91 status: String,
92 attempts: u32,
93 response_status: Option<u16>,
94 response_body: Option<String>,
95 error: Option<String>,
96 duration_ms: Option<u32>,
97 created_at: String,
98 delivered_at: Option<String>,
99 next_attempt_at: Option<String>,
100}
101
102impl From<DeliveryRow> for HookDelivery {
103 fn from(row: DeliveryRow) -> Self {
104 HookDelivery {
105 id: row.id,
106 hook_id: row.hook_id,
107 event_id: row.event_id,
108 event: row.event,
109 status: row.status,
110 attempts: row.attempts,
111 response_status: row.response_status,
112 response_body: row.response_body,
113 error: row.error,
114 duration_ms: row.duration_ms,
115 payload: row.payload,
116 created_at: row.created_at,
117 delivered_at: row.delivered_at,
118 next_attempt_at: row.next_attempt_at,
119 }
120 }
121}
122
123#[derive(Deserialize)]
124struct NameRow {
125 namespace: String,
126 name: String,
127}
128
129/// What one attempt to send came to.
130struct Attempt {
131 status: Option<u16>,
132 body: Option<String>,
133 error: Option<String>,
134 duration_ms: u32,
135}
136
137impl Attempt {
138 fn delivered(&self) -> bool {
139 self.status.is_some_and(|status| (200..300).contains(&status))
140 }
141}
142
143fn optional(value: Option<&str>) -> JsValue {
144 value.map_or(JsValue::NULL, JsValue::from)
145}
146
147fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
148 Outcome::fail(code, message)
149}
150
151/// The workspace itself, as the one asking: it can see its own
152/// repositories, private ones too, and no one else's.
153fn workspace_viewer(slug: &str) -> Viewer {
154 Some(User {
155 id: String::new(),
156 username: slug.to_owned(),
157 kind: PrincipalKind::Workspace,
158 verified: true,
159 workspaces: vec![Membership::member(slug.to_owned())],
160 ..User::default()
161 })
162}
163
164struct Webhooks {
165 db: D1Database,
166 sealer: Option<Sealer>,
167 repos: Fetcher,
168 identity: Fetcher,
169}
170
171impl Webhooks {
172 fn new(env: &Env) -> Result<Self> {
173 Ok(Webhooks {
174 db: env.d1("DB")?,
175 sealer: env.secret("WEBHOOKS_KEY").ok().and_then(|key| Sealer::new(&key.to_string())),
176 repos: env.service("REPOS")?,
177 identity: env.service("IDENTITY")?,
178 })
179 }
180
181 // --- Who may do what --------------------------------------------------------
182
183 /// The repository a repository's webhooks are for, if it is the
184 /// workspace's and the viewer can see it.
185 async fn repository(&self, owner: &HookOwner, viewer: &Viewer) -> Result<Option<Repo>> {
186 let Some(path) = &owner.repo else {
187 return Ok(None);
188 };
189 let found: Outcome<Repo> = g1t_kit::call(
190 &self.repos,
191 "get",
192 &GetArgs {
193 path: path.clone(),
194 viewer: viewer.clone(),
195 },
196 )
197 .await?;
198 Ok(found
199 .into_result()
200 .ok()
201 .filter(|repo| repo.namespace == owner.workspace && repo.fork_of.is_none()))
202 }
203
204 /// Whether `viewer` may see (or, `managing`, change) `owner`'s webhooks,
205 /// and the repository when they are a repository's. A repository's
206 /// need the Admin role on it, to see as to change, since they carry
207 /// its events out; the workspace's are its members' to see and its
208 /// owners' to change. An agent's token changes neither.
209 async fn allowed(&self, viewer: &Viewer, owner: &HookOwner, managing: bool) -> Result<Outcome<Option<Repo>>> {
210 let Some(user) = viewer.as_ref() else {
211 return Ok(fail(FailureCode::Forbidden, "Sign in to see webhooks."));
212 };
213 if managing && user.kind == PrincipalKind::Agent {
214 return Ok(fail(FailureCode::Forbidden, "An agent cannot manage webhooks."));
215 }
216 if owner.repo.is_some() {
217 let Some(repo) = self.repository(owner, viewer).await? else {
218 return Ok(fail(FailureCode::NotFound, "There is no such repository in this workspace."));
219 };
220 if !access::can(viewer.as_ref(), &repo, Capability::ManageIntegrations) {
221 return Ok(fail(
222 FailureCode::Forbidden,
223 access::needs(Capability::ManageIntegrations, &format!("{}/{}", repo.namespace, repo.name)),
224 ));
225 }
226 return Ok(Outcome::Ok(Some(repo)));
227 }
228 if !user.is_member(&owner.workspace) {
229 return Ok(fail(FailureCode::Forbidden, format!("Only members of {} can see its webhooks.", owner.workspace)));
230 }
231 if managing && user.role_in(&owner.workspace) != Some(Role::Owner) {
232 return Ok(fail(FailureCode::Forbidden, "Only an owner can manage a workspace's own webhooks."));
233 }
234 Ok(Outcome::Ok(None))
235 }
236
237 fn owner(mut owner: HookOwner) -> HookOwner {
238 owner.workspace = owner.workspace.to_lowercase();
239 if let Some(repo) = &mut owner.repo {
240 repo.namespace = repo.namespace.to_lowercase();
241 }
242 owner
243 }
244
245 /// The webhook, if it belongs to `owner`.
246 async fn hook_of(&self, owner: &HookOwner, id: &str) -> Result<Option<HookRow>> {
247 let row = self
248 .db
249 .prepare("SELECT * FROM hooks WHERE id = ? AND workspace = ?")
250 .bind(&[id.into(), owner.workspace.as_str().into()])?
251 .first::<HookRow>(None)
252 .await?;
253 let repo = owner.repo.as_ref().map(|path| format!("{}/{}", path.namespace, path.name));
254 Ok(row.filter(|row| match &repo {
255 Some(repo) => row.scope == "repo" && row.repo.as_deref().is_some_and(|r| r.eq_ignore_ascii_case(repo)),
256 None => row.scope == "workspace",
257 }))
258 }
259
260 // --- Managing webhooks ------------------------------------------------------
261
262 async fn list(&self, a: ListArgs) -> Result<Outcome<Vec<Hook>>> {
263 let owner = Self::owner(a.owner);
264 if let Outcome::Fail(refused) = self.allowed(&a.viewer, &owner, false).await? {
265 return Ok(Outcome::Fail(refused));
266 }
267 let rows = match &owner.repo {
268 Some(path) => self
269 .db
270 .prepare("SELECT * FROM hooks WHERE scope = 'repo' AND workspace = ? AND lower(repo) = lower(?) ORDER BY id")
271 .bind(&[owner.workspace.as_str().into(), format!("{}/{}", path.namespace, path.name).into()])?,
272 None => self
273 .db
274 .prepare("SELECT * FROM hooks WHERE scope = 'workspace' AND workspace = ? ORDER BY id")
275 .bind(&[owner.workspace.as_str().into()])?,
276 }
277 .all()
278 .await?
279 .results::<HookRow>()?;
280 Ok(Outcome::Ok(rows.iter().map(HookRow::to_hook).collect()))
281 }
282
283 async fn create(&self, a: CreateArgs) -> Result<Outcome<CreatedHook>> {
284 let owner = Self::owner(a.owner);
285 let repo = match self.allowed(&Some(a.actor.clone()), &owner, true).await? {
286 Outcome::Ok(repo) => repo,
287 Outcome::Fail(refused) => return Ok(Outcome::Fail(refused)),
288 };
289 let Some(sealer) = &self.sealer else {
290 return Ok(fail(FailureCode::Conflict, "Webhooks are not set up on this g1t: it has no key to keep secrets with."));
291 };
292 let url = a.url.trim().to_owned();
293 if let Err(problem) = deliver::check_url(&url) {
294 return Ok(fail(FailureCode::Invalid, problem));
295 }
296 let events = match deliver::tidy_events(&a.events) {
297 Ok(events) => events,
298 Err(problem) => return Ok(fail(FailureCode::Invalid, problem)),
299 };
300 let given = a.secret.map(|secret| secret.trim().to_owned()).filter(|secret| !secret.is_empty());
301 let made = given.is_none();
302 let secret = given.unwrap_or_else(|| format!("whsec_{}", g1t_secrets::random_hex(24)));
303 let now = now_ms();
304 let id = new_id("hk", now);
305 self.db
306 .prepare(
307 "INSERT INTO hooks (id, scope, workspace, repo_id, repo, url, events, secret, secret_hint, created_by, created_at)
308 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
309 )
310 .bind(&[
311 id.as_str().into(),
312 if repo.is_some() { "repo" } else { "workspace" }.into(),
313 owner.workspace.as_str().into(),
314 optional(repo.as_ref().map(|repo| repo.id.as_str())),
315 optional(repo.as_ref().map(|repo| format!("{}/{}", repo.namespace, repo.name)).as_deref()),
316 url.as_str().into(),
317 serde_json::to_string(&events)?.into(),
318 sealer.seal(&secret, &id).into(),
319 g1t_secrets::hint(&secret).into(),
320 a.actor.username.as_str().into(),
321 rfc3339(now).into(),
322 ])?
323 .run()
324 .await?;
325 let Some(row) = self.hook_of(&owner, &id).await? else {
326 return Ok(fail(FailureCode::NotFound, "The webhook was not saved."));
327 };
328 // Tells the receiver it is wired up, and shows whether it answers.
329 self.send_ping(&row).await?;
330 let row = self.hook_of(&owner, &id).await?.unwrap_or(row);
331 Ok(Outcome::Ok(CreatedHook {
332 hook: row.to_hook(),
333 secret: made.then_some(secret),
334 }))
335 }
336
337 async fn update(&self, a: UpdateArgs) -> Result<Outcome<Hook>> {
338 let owner = Self::owner(a.owner);
339 if let Outcome::Fail(refused) = self.allowed(&Some(a.actor.clone()), &owner, true).await? {
340 return Ok(Outcome::Fail(refused));
341 }
342 let Some(row) = self.hook_of(&owner, &a.id).await? else {
343 return Ok(fail(FailureCode::NotFound, "No such webhook."));
344 };
345 let url = a.url.map(|url| url.trim().to_owned()).unwrap_or(row.url.clone());
346 if let Err(problem) = deliver::check_url(&url) {
347 return Ok(fail(FailureCode::Invalid, problem));
348 }
349 let events = match a.events {
350 Some(events) => match deliver::tidy_events(&events) {
351 Ok(events) => events,
352 Err(problem) => return Ok(fail(FailureCode::Invalid, problem)),
353 },
354 None => row.events(),
355 };
356 let active = a.active.unwrap_or(row.active != 0);
357 self.db
358 .prepare("UPDATE hooks SET url = ?, events = ?, active = ? WHERE id = ?")
359 .bind(&[url.into(), serde_json::to_string(&events)?.into(), (active as u32).into(), row.id.as_str().into()])?
360 .run()
361 .await?;
362 Ok(match self.hook_of(&owner, &row.id).await? {
363 Some(row) => Outcome::Ok(row.to_hook()),
364 None => fail(FailureCode::NotFound, "No such webhook."),
365 })
366 }
367
368 async fn delete(&self, a: HookArgs) -> Result<Outcome<bool>> {
369 let owner = Self::owner(a.owner);
370 if let Outcome::Fail(refused) = self.allowed(&Some(a.actor.clone()), &owner, true).await? {
371 return Ok(Outcome::Fail(refused));
372 }
373 let Some(row) = self.hook_of(&owner, &a.id).await? else {
374 return Ok(fail(FailureCode::NotFound, "No such webhook."));
375 };
376 self.db
377 .batch(vec![
378 self.db.prepare("DELETE FROM hooks WHERE id = ?").bind(&[row.id.as_str().into()])?,
379 self.db.prepare("DELETE FROM deliveries WHERE hook_id = ?").bind(&[row.id.as_str().into()])?,
380 ])
381 .await?;
382 Ok(Outcome::Ok(true))
383 }
384
385 async fn ping(&self, a: HookArgs) -> Result<Outcome<HookDelivery>> {
386 let owner = Self::owner(a.owner);
387 if let Outcome::Fail(refused) = self.allowed(&Some(a.actor.clone()), &owner, true).await? {
388 return Ok(Outcome::Fail(refused));
389 }
390 let Some(row) = self.hook_of(&owner, &a.id).await? else {
391 return Ok(fail(FailureCode::NotFound, "No such webhook."));
392 };
393 let id = self.send_ping(&row).await?;
394 Ok(self.delivery(&id).await?.map_or_else(|| fail(FailureCode::NotFound, "The ping was not recorded."), Outcome::Ok))
395 }
396
397 async fn deliveries(&self, a: DeliveriesArgs) -> Result<Outcome<Vec<HookDelivery>>> {
398 let owner = Self::owner(a.owner);
399 if let Outcome::Fail(refused) = self.allowed(&a.viewer, &owner, false).await? {
400 return Ok(Outcome::Fail(refused));
401 }
402 if self.hook_of(&owner, &a.id).await?.is_none() {
403 return Ok(fail(FailureCode::NotFound, "No such webhook."));
404 }
405 let rows = self
406 .db
407 .prepare("SELECT * FROM deliveries WHERE hook_id = ? ORDER BY id DESC LIMIT ?")
408 .bind(&[a.id.as_str().into(), DELIVERIES_SHOWN.into()])?
409 .all()
410 .await?
411 .results::<DeliveryRow>()?;
412 Ok(Outcome::Ok(rows.into_iter().map(HookDelivery::from).collect()))
413 }
414
415 async fn redeliver(&self, a: RedeliverArgs) -> Result<Outcome<HookDelivery>> {
416 let owner = Self::owner(a.owner);
417 if let Outcome::Fail(refused) = self.allowed(&Some(a.actor.clone()), &owner, true).await? {
418 return Ok(Outcome::Fail(refused));
419 }
420 let Some(original) = self.delivery(&a.delivery_id).await? else {
421 return Ok(fail(FailureCode::NotFound, "No such delivery."));
422 };
423 let Some(row) = self.hook_of(&owner, &original.hook_id).await? else {
424 return Ok(fail(FailureCode::NotFound, "No such delivery."));
425 };
426 // A new delivery of the same payload: its own attempts and log.
427 let id = self.enqueue(&row, "", &original.event, &original.payload).await?;
428 if let Some(id) = &id {
429 self.attempt(&row, id).await?;
430 }
431 Ok(match id {
432 Some(id) => self.delivery(&id).await?.map_or_else(|| fail(FailureCode::NotFound, "No such delivery."), Outcome::Ok),
433 None => fail(FailureCode::Conflict, "That delivery could not be made again."),
434 })
435 }
436
437 async fn delivery(&self, id: &str) -> Result<Option<HookDelivery>> {
438 Ok(self
439 .db
440 .prepare("SELECT * FROM deliveries WHERE id = ?")
441 .bind(&[id.into()])?
442 .first::<DeliveryRow>(None)
443 .await?
444 .map(HookDelivery::from))
445 }
446
447 // --- Delivering -------------------------------------------------------------
448
449 /// Records a delivery to make. `None` when this event was already
450 /// delivered to this webhook.
451 async fn enqueue(&self, hook: &HookRow, event_id: &str, event: &str, payload: &str) -> Result<Option<String>> {
452 let now = now_ms();
453 let id = new_id("dlv", now);
454 let inserted = self
455 .db
456 .prepare(
457 "INSERT OR IGNORE INTO deliveries (id, hook_id, event_id, event, payload, status, created_at, next_attempt_at)
458 VALUES (?, ?, ?, ?, ?, 'pending', ?, ?) RETURNING id",
459 )
460 .bind(&[
461 id.as_str().into(),
462 hook.id.as_str().into(),
463 event_id.into(),
464 event.into(),
465 payload.into(),
466 rfc3339(now).into(),
467 rfc3339(now).into(),
468 ])?
469 .first::<Value>(None)
470 .await?;
471 Ok(inserted.map(|_| id))
472 }
473
474 async fn send_ping(&self, hook: &HookRow) -> Result<String> {
475 let body = deliver::ping(&hook.id, &hook.url, &hook.events(), &rfc3339(now_ms())).to_string();
476 let id = self.enqueue(hook, "", "ping", &body).await?.unwrap_or_default();
477 self.attempt(hook, &id).await?;
478 Ok(id)
479 }
480
481 /// Sends one delivery once, and records how it went.
482 async fn attempt(&self, hook: &HookRow, delivery_id: &str) -> Result<()> {
483 let Some(delivery) = self.delivery(delivery_id).await? else {
484 return Ok(());
485 };
486 let Some(secret) = self.sealer.as_ref().and_then(|sealer| sealer.open(&hook.secret, &hook.id)) else {
487 return Ok(());
488 };
489 let attempt = send(&hook.url, &hook.id, &delivery, &secret).await;
490 let attempts = delivery.attempts + 1;
491 let now = now_ms();
492 let (status, next) = if attempt.delivered() {
493 ("delivered", None)
494 } else {
495 match deliver::retry_after(attempts) {
496 Some(wait) => ("pending", Some(rfc3339(now + wait * 1000))),
497 None => ("failed", None),
498 }
499 };
500 self.db
501 .batch(vec![
502 self.db
503 .prepare(
504 "UPDATE deliveries SET status = ?, attempts = ?, response_status = ?, response_body = ?, error = ?,
505 duration_ms = ?, delivered_at = ?, next_attempt_at = ? WHERE id = ?",
506 )
507 .bind(&[
508 status.into(),
509 attempts.into(),
510 attempt.status.map_or(JsValue::NULL, |status| status.into()),
511 optional(attempt.body.as_deref()),
512 optional(attempt.error.as_deref()),
513 attempt.duration_ms.into(),
514 optional(attempt.delivered().then(|| rfc3339(now)).as_deref()),
515 optional(next.as_deref()),
516 delivery_id.into(),
517 ])?,
518 self.db
519 .prepare("UPDATE hooks SET last_status = ?, last_delivered_at = ? WHERE id = ?")
520 .bind(&[status.into(), rfc3339(now).into(), hook.id.as_str().into()])?,
521 ])
522 .await?;
523 Ok(())
524 }
525
526 /// Which workspace a repository is in, and its name.
527 async fn repo_name(&self, repo_id: &str, workspaces: &[String]) -> Result<Option<NameRow>> {
528 if let Some(known) = self
529 .db
530 .prepare("SELECT namespace, name FROM repo_names WHERE repo_id = ?")
531 .bind(&[repo_id.into()])?
532 .first::<NameRow>(None)
533 .await?
534 {
535 return Ok(Some(known));
536 }
537 // Asked as each workspace with webhooks in turn: each sees only its
538 // own private repositories.
539 for workspace in workspaces {
540 let found: Outcome<Repo> = g1t_kit::call(
541 &self.repos,
542 "get_by_id",
543 &GetByIdArgs {
544 id: repo_id.to_owned(),
545 viewer: workspace_viewer(workspace),
546 },
547 )
548 .await?;
549 if let Outcome::Ok(repo) = found
550 && repo.fork_of.is_none()
551 {
552 self.remember(repo_id, &repo.namespace, &repo.name).await?;
553 return Ok(Some(NameRow {
554 namespace: repo.namespace,
555 name: repo.name,
556 }));
557 }
558 }
559 Ok(None)
560 }
561
562 async fn remember(&self, repo_id: &str, namespace: &str, name: &str) -> Result<()> {
563 self.db
564 .prepare("INSERT OR REPLACE INTO repo_names (repo_id, namespace, name) VALUES (?, ?, ?)")
565 .bind(&[repo_id.into(), namespace.into(), name.into()])?
566 .run()
567 .await?;
568 Ok(())
569 }
570
571 /// An event from the bus, delivered to every webhook that wants it.
572 async fn on_event(&self, event: &Event) -> Result<()> {
573 if event.kind == "repo.created"
574 && let (Some(id), Some(namespace), Some(name)) =
575 (event.data["repoId"].as_str(), event.data["namespace"].as_str(), event.data["name"].as_str())
576 {
577 self.remember(id, namespace, name).await?;
578 }
579 let Some(repo_id) = event.repo_id.as_deref().or_else(|| event.data["repoId"].as_str()) else {
580 return Ok(());
581 };
582 let mut hooks = self
583 .db
584 .prepare("SELECT * FROM hooks WHERE active = 1 AND scope = 'repo' AND repo_id = ?")
585 .bind(&[repo_id.into()])?
586 .all()
587 .await?
588 .results::<HookRow>()?;
589 let workspaces: Vec<String> = self
590 .db
591 .prepare("SELECT DISTINCT workspace FROM hooks WHERE active = 1 AND scope = 'workspace'")
592 .all()
593 .await?
594 .results::<Value>()?
595 .into_iter()
596 .filter_map(|row| row["workspace"].as_str().map(str::to_owned))
597 .collect();
598 let name = if hooks.is_empty() && workspaces.is_empty() {
599 None
600 } else {
601 // A deleted repository looks missing to repos; its own events
602 // name it.
603 self.repo_name(repo_id, &workspaces)
604 .await?
605 .or_else(|| deliver::named_in(event).map(|(namespace, name)| NameRow { namespace, name }))
606 };
607 if let Some(name) = &name {
608 hooks.extend(
609 self.db
610 .prepare("SELECT * FROM hooks WHERE active = 1 AND scope = 'workspace' AND workspace = ?")
611 .bind(&[name.namespace.as_str().into()])?
612 .all()
613 .await?
614 .results::<HookRow>()?,
615 );
616 }
617 let wanted: Vec<&HookRow> = hooks.iter().filter(|hook| deliver::wants(&hook.events(), &event.kind)).collect();
618 if wanted.is_empty() {
619 return Ok(());
620 }
621 // Who caused it, by name: people and workspaces, or g1t's agent.
622 let actor_name = match &event.actor {
623 Some(id) => {
624 let names: std::collections::HashMap<String, String> =
625 g1t_kit::call(&self.identity, "usernames", &UsernamesArgs { ids: vec![id.clone()] }).await?;
626 names.get(id).cloned().or_else(|| (id == AGENT_ID).then(|| AGENT_NAME.to_owned()))
627 }
628 None => None,
629 };
630 for hook in wanted {
631 let full_name = hook
632 .repo
633 .clone()
634 .or_else(|| name.as_ref().map(|name| format!("{}/{}", name.namespace, name.name)))
635 .unwrap_or_default();
636 let payload = deliver::payload(event, &hook.workspace, Some((repo_id, &full_name)), actor_name.as_deref()).to_string();
637 if let Some(id) = self.enqueue(hook, &event.id, &event.kind, &payload).await? {
638 self.attempt(hook, &id).await?;
639 }
640 }
641 Ok(())
642 }
643
644 /// Tries again what is due, and forgets what is old.
645 async fn sweep(&self) -> Result<()> {
646 let now = now_ms();
647 let due = self
648 .db
649 .prepare("SELECT * FROM deliveries WHERE status = 'pending' AND next_attempt_at <= ? ORDER BY next_attempt_at LIMIT ?")
650 .bind(&[rfc3339(now).into(), SWEEP.into()])?
651 .all()
652 .await?
653 .results::<DeliveryRow>()?;
654 for delivery in due {
655 let hook = self
656 .db
657 .prepare("SELECT * FROM hooks WHERE id = ?")
658 .bind(&[delivery.hook_id.as_str().into()])?
659 .first::<HookRow>(None)
660 .await?;
661 match hook {
662 Some(hook) if hook.active != 0 => self.attempt(&hook, &delivery.id).await?,
663 // A webhook turned off or removed takes its retries with it.
664 _ => {
665 self.db
666 .prepare("UPDATE deliveries SET status = 'failed', next_attempt_at = NULL WHERE id = ?")
667 .bind(&[delivery.id.as_str().into()])?
668 .run()
669 .await?;
670 }
671 }
672 }
673 self.db
674 .prepare("DELETE FROM deliveries WHERE created_at < ?")
675 .bind(&[rfc3339(now.saturating_sub(KEPT_DAYS * 24 * 60 * 60 * 1000)).into()])?
676 .run()
677 .await?;
678 Ok(())
679 }
680}
681
682/// One HTTPS POST of a delivery's payload, signed, given ten seconds.
683async fn send(url: &str, hook_id: &str, delivery: &HookDelivery, secret: &str) -> Attempt {
684 let started = now_ms();
685 let elapsed = || (now_ms() - started) as u32;
686 let request = (|| -> Result<Request> {
687 let headers = Headers::new();
688 headers.set("content-type", "application/json")?;
689 headers.set("user-agent", "g1t-webhooks/1 (+https://docs.g1t.sh/guides/webhooks/)")?;
690 headers.set("x-g1t-event", &delivery.event)?;
691 headers.set("x-g1t-delivery", &delivery.id)?;
692 headers.set("x-g1t-hook", hook_id)?;
693 headers.set("x-g1t-signature-256", &deliver::signature(secret, &delivery.payload))?;
694 let mut init = RequestInit::new();
695 init.with_method(Method::Post).with_headers(headers).with_body(Some(delivery.payload.clone().into()));
696 Request::new_with_init(url, &init)
697 })();
698 let request = match request {
699 Ok(request) => request,
700 Err(error) => {
701 return Attempt {
702 status: None,
703 body: None,
704 error: Some(format!("The request could not be made: {error}")),
705 duration_ms: 0,
706 };
707 }
708 };
709 let fetcher = Fetch::Request(request);
710 let fetch = Box::pin(fetcher.send());
711 let timeout = Box::pin(Delay::from(TIMEOUT));
712 match select(fetch, timeout).await {
713 Either::Left((Ok(mut response), _)) => {
714 let status = response.status_code();
715 let body: String = response.text().await.unwrap_or_default().chars().take(RESPONSE_KEPT).collect();
716 Attempt {
717 status: Some(status),
718 body: Some(body),
719 error: None,
720 duration_ms: elapsed(),
721 }
722 }
723 Either::Left((Err(error), _)) => Attempt {
724 status: None,
725 body: None,
726 error: Some(format!("The receiver could not be reached: {error}")),
727 duration_ms: elapsed(),
728 },
729 Either::Right(_) => Attempt {
730 status: None,
731 body: None,
732 error: Some(format!("The receiver did not answer within {} seconds.", TIMEOUT.as_secs())),
733 duration_ms: elapsed(),
734 },
735 }
736}
737
738#[event(fetch)]
739async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
740 let Some(method) = rpc_method(&request) else {
741 return Response::error("Not found", 404);
742 };
743 let body: Value = request.json().await?;
744 let service = Webhooks::new(&env)?;
745 match method.as_str() {
746 "list" => reply(&service.list(args(body)?).await?),
747 "create" => reply(&service.create(args(body)?).await?),
748 "update" => reply(&service.update(args(body)?).await?),
749 "delete" => reply(&service.delete(args(body)?).await?),
750 "ping" => reply(&service.ping(args(body)?).await?),
751 "deliveries" => reply(&service.deliveries(args(body)?).await?),
752 "redeliver" => reply(&service.redeliver(args(body)?).await?),
753 _ => Response::error("Unknown method", 404),
754 }
755}
756
757/// Events from the bus, on this service's own queue.
758#[event(queue)]
759async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
760 let service = Webhooks::new(&env)?;
761 for message in batch.messages()? {
762 // A workspace renamed: its rows move to the slug it has now.
763 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), rename::STATEMENTS).await? {
764 message.ack();
765 continue;
766 }
767 // A repository renamed or transferred: its rows follow its new path,
768 // and then the event is delivered like any other, naming it there.
769 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), rename::TRANSFERRED).await? {
770 service.on_event(message.body()).await?;
771 message.ack();
772 continue;
773 }
774 // A repository purged: its last event is delivered, then its own
775 // webhooks go.
776 if message.body().kind == "repo.purged" {
777 service.on_event(message.body()).await?;
778 g1t_kit::lifecycle::on_purged(&env.d1("DB")?, message.body(), rename::PURGED).await?;
779 message.ack();
780 continue;
781 }
782 // A workspace deleted: what it kept for itself goes.
783 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), rename::DELETED).await? {
784 message.ack();
785 continue;
786 }
787 service.on_event(message.body()).await?;
788 message.ack();
789 }
790 Ok(())
791}
792
793/// Every minute: retries that are due, and deliveries old enough to forget.
794#[event(scheduled)]
795async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
796 match Webhooks::new(&env) {
797 Ok(service) => {
798 if let Err(error) = service.sweep().await {
799 worker::console_error!("webhooks: the sweep failed: {error}");
800 }
801 }
802 Err(error) => worker::console_error!("webhooks: could not start: {error}"),
803 }
804}
805