Skip to content
881 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

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

This file's history is long; its oldest lines are credited to the oldest commit read.