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

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