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

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