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

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