Skip to content
507 linesCodeBlameRaw
1//! The events service: the bus every state change in g1t is published on,
2//! and its durable log.
3//!
4//! Publishing puts events on a queue and returns. The queue consumer writes
5//! them to the log and passes each batch on to each subscriber's own
6//! queue, so a slow or failing subscriber holds up nobody else. Each
7//! subscriber is sent only the types it acts on
8//! (`g1t_contracts::subscribers`), in batches the queues take (fanout.rs).
9//!
10//! It keeps two more things beside the log: the audit log (audit.rs) and
11//! each person's inbox (inbox.rs), written as events arrive, with who
12//! follows what (subscriptions.rs).
13//!
14//! Other services reach it over `POST /rpc/<method>`; see
15//! `g1t_contracts::events`, `audit` and `inbox` for the methods and their
16//! arguments.
17
18mod audit;
19mod fanout;
20mod inbox;
21mod subscriptions;
22
23use g1t_contracts::events::{Event, ListArgs, PublishArgs};
24use g1t_contracts::new_id;
25use g1t_contracts::time::rfc3339;
26use g1t_kit::{args, js, now_ms, reply, rpc_method};
27use serde::Deserialize;
28use std::collections::{HashMap, HashSet};
29use worker::js_sys::{Array, Object};
30use worker::wasm_bindgen::{JsCast, JsValue};
31use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event};
32
33const DEFAULT_PAGE: u32 = 50;
34const MAX_PAGE: u32 = 200;
35/// Statuses and check runs reported on commits: delivered to webhooks, and
36/// left out of a timeline unless asked for by type.
37const REPORTING: [&str; 7] = [
38 "status.created",
39 "check_run.created",
40 "check_run.completed",
41 "check_run.rerequested",
42 "check_run.requested_action",
43 "check_suite.completed",
44 "check_suite.rerequested",
45];
46/// Every binding whose name starts with this is a queue that receives
47/// events: one per subscribing service, sent the types it routes.
48const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_";
49/// How long a record of which queues a batch reached is kept for its retries.
50const FANOUT_KEEP_MS: u64 = 24 * 60 * 60 * 1000;
51
52#[derive(Deserialize)]
53struct EventRow {
54 id: String,
55 #[serde(rename = "type")]
56 kind: String,
57 source: String,
58 time: String,
59 repo_id: Option<String>,
60 actor: Option<String>,
61 /// JSON.
62 data: String,
63}
64
65impl From<EventRow> for Event {
66 fn from(row: EventRow) -> Self {
67 Event {
68 id: row.id,
69 kind: row.kind,
70 source: row.source,
71 time: row.time,
72 repo_id: row.repo_id,
73 actor: row.actor,
74 data: serde_json::from_str(&row.data).unwrap_or_default(),
75 }
76 }
77}
78
79fn optional(value: &Option<String>) -> JsValue {
80 value.as_deref().map_or(JsValue::NULL, JsValue::from)
81}
82
83#[derive(Deserialize)]
84struct FanoutRow {
85 event_id: String,
86 /// JSON array of binding names.
87 bindings: String,
88}
89
90/// Sends `events` to a queue binding, one message each, in as many
91/// `sendBatch` calls as the queue's limits need (fanout.rs). Events are
92/// expected to have been fit to a message.
93async fn send(queue: &JsValue, events: &[&Event]) -> Result<()> {
94 for chunk in fanout::chunks(events) {
95 let messages = Array::new();
96 for event in chunk {
97 let message = Object::new();
98 js::set(&message, "body", &js::to_js(event)?);
99 messages.push(&message);
100 }
101 js::call(queue, "sendBatch", &[messages.into()]).await?;
102 }
103 Ok(())
104}
105
106struct Events {
107 db: D1Database,
108 env: Env,
109}
110
111impl Events {
112 /// Assigns each event its id and time and puts it on the bus.
113 async fn publish(&self, a: PublishArgs) -> Result<()> {
114 if a.events.is_empty() {
115 return Ok(());
116 }
117 let now = now_ms();
118 let events: Vec<Event> = a
119 .events
120 .into_iter()
121 .map(|event| {
122 let mut event = Event {
123 id: new_id("evt", now),
124 kind: event.kind,
125 source: event.source,
126 time: rfc3339(now),
127 repo_id: event.repo_id,
128 actor: event.actor,
129 data: event.data,
130 };
131 // Too large for one message: its long text is shortened.
132 fanout::fit(&mut event);
133 event
134 })
135 .collect();
136 let events: Vec<&Event> = events.iter().collect();
137 send(&js::binding(&self.env, "BUS")?, &events).await
138 }
139
140 /// Newest first. Callers must have checked that the viewer may see the
141 /// repository asked about.
142 async fn list(&self, a: ListArgs) -> Result<Vec<Event>> {
143 let (filter, binds) = list_filter(&a);
144 let mut values: Vec<JsValue> = binds
145 .into_iter()
146 .map(|bind| match bind {
147 Bind::Text(text) => JsValue::from(text),
148 Bind::Number(number) => JsValue::from(number),
149 })
150 .collect();
151 values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into());
152 let rows = self
153 .db
154 .prepare(format!(
155 "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?"
156 ))
157 .bind(&values)?
158 .all()
159 .await?
160 .results::<EventRow>()?;
161 Ok(rows.into_iter().map(Event::from).collect())
162 }
163
164 /// Writes a batch from the bus to the log, notes who now follows what
165 /// (inbox.rs), hands each subscriber the events it routes, then tells
166 /// the people it concerns.
167 ///
168 /// Everything before the hand-off is safe to repeat. A retried batch
169 /// skips the queues that already have it: when a hand-off fails part
170 /// way, which queues each event reached is written down first.
171 async fn deliver(&self, events: &[Event]) -> Result<()> {
172 let mut statements = Vec::with_capacity(events.len() + 1);
173 for event in events {
174 statements.push(
175 self.db
176 .prepare(
177 // Redelivered batches must not duplicate log rows.
178 "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data)
179 VALUES (?, ?, ?, ?, ?, ?, ?)",
180 )
181 .bind(&[
182 event.id.as_str().into(),
183 event.kind.as_str().into(),
184 event.source.as_str().into(),
185 event.time.as_str().into(),
186 optional(&event.repo_id),
187 optional(&event.actor),
188 serde_json::to_string(&event.data)?.into(),
189 ])?,
190 );
191 }
192 let ids: Vec<&str> = events.iter().map(|event| event.id.as_str()).collect();
193 // Read in the same round trip: nearly always nothing.
194 statements.push(
195 self.db
196 .prepare("SELECT event_id, bindings FROM fanout_sent WHERE event_id IN (SELECT value FROM json_each(?))")
197 .bind(&[serde_json::to_string(&ids)?.into()])?,
198 );
199 let results = self.db.batch(statements).await?;
200 let mut sent: HashMap<String, HashSet<String>> = HashMap::new();
201 if let Some(rows) = results.last() {
202 for row in rows.results::<FanoutRow>()? {
203 let bindings: Vec<String> = serde_json::from_str(&row.bindings).unwrap_or_default();
204 sent.insert(row.event_id, bindings.into_iter().collect());
205 }
206 }
207 audit::follow_renames(&self.db, events).await?;
208 // Before the hand-off, so its failing does not send the batch to
209 // every queue again.
210 inbox::follow(&self.db, events).await?;
211
212 let bindings: &JsValue = self.env.as_ref();
213 let mut reached: Vec<String> = Vec::new();
214 let mut failed = None;
215 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
216 let Some(name) = name
217 .as_string()
218 .filter(|name| name.starts_with(SUBSCRIBER_PREFIX))
219 else {
220 continue;
221 };
222 let routed: Vec<&Event> = events
223 .iter()
224 .filter(|event| g1t_contracts::subscribers::routed(&name, &event.kind))
225 .filter(|event| !sent.get(&event.id).is_some_and(|reached| reached.contains(&name)))
226 .collect();
227 if routed.is_empty() {
228 continue;
229 }
230 match send(&js::binding(&self.env, &name)?, &routed).await {
231 Ok(()) => reached.push(name),
232 Err(error) => {
233 worker::console_error!("events: passing {} events to {name} failed: {error}", routed.len());
234 failed = Some(error);
235 }
236 }
237 }
238 if let Some(error) = failed {
239 self.note_reached(events, &sent, &reached).await?;
240 return Err(error);
241 }
242 let (work, repos, identity) = (
243 self.env.service("WORK")?,
244 self.env.service("REPOS")?,
245 self.env.service("IDENTITY")?,
246 );
247 let sources = inbox::Sources {
248 work: &work,
249 repos: &repos,
250 identity: &identity,
251 };
252 inbox::deliver(&self.db, &sources, events).await;
253 Ok(())
254 }
255
256 /// Writes down which queues each event of a batch has reached, with
257 /// those it had reached before, for the batch's retry.
258 async fn note_reached(
259 &self,
260 events: &[Event],
261 sent: &HashMap<String, HashSet<String>>,
262 reached: &[String],
263 ) -> Result<()> {
264 let now = rfc3339(now_ms());
265 let mut statements = Vec::with_capacity(events.len());
266 for event in events {
267 let mut bindings: Vec<&str> = reached.iter().map(String::as_str).collect();
268 if let Some(before) = sent.get(&event.id) {
269 bindings.extend(before.iter().map(String::as_str));
270 }
271 bindings.sort_unstable();
272 bindings.dedup();
273 statements.push(
274 self.db
275 .prepare(
276 "INSERT INTO fanout_sent (event_id, bindings, created_at) VALUES (?, ?, ?)
277 ON CONFLICT (event_id) DO UPDATE SET bindings = excluded.bindings",
278 )
279 .bind(&[event.id.as_str().into(), serde_json::to_string(&bindings)?.into(), now.as_str().into()])?,
280 );
281 }
282 self.db.batch(statements).await?;
283 Ok(())
284 }
285}
286
287#[event(fetch)]
288async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
289 let Some(method) = rpc_method(&request) else {
290 return Response::error("Not found", 404);
291 };
292 let body: serde_json::Value = request.json().await?;
293 let events = Events {
294 db: env.d1("DB")?,
295 env,
296 };
297 match method.as_str() {
298 "publish" => reply(&events.publish(args(body)?).await?),
299 "list" => reply(&events.list(args(body)?).await?),
300 "audit_record" => reply(&audit::record(&events.db, args(body)?).await?),
301 "audit_list" => reply(&audit::list(&events.db, args(body)?).await?),
302 "inbox_list" => {
303 let repos = events.env.service("REPOS")?;
304 reply(&inbox::list(&events.db, &repos, args(body)?).await?)
305 }
306 "inbox_counts" => reply(&inbox::counts(&events.db, args(body)?).await?),
307 "inbox_mark" => reply(&inbox::mark(&events.db, args(body)?).await?),
308 "inbox_thread" => {
309 let (repos, work) = (events.env.service("REPOS")?, events.env.service("WORK")?);
310 reply(&inbox::thread(&events.db, &repos, &work, args(body)?).await?)
311 }
312 "inbox_subscription" => {
313 let work = events.env.service("WORK")?;
314 reply(&subscriptions::subscription(&events.db, &work, args(body)?).await?)
315 }
316 "inbox_subscribe" => {
317 let work = events.env.service("WORK")?;
318 reply(&subscriptions::subscribe(&events.db, &work, args(body)?).await?)
319 }
320 "inbox_watching" => reply(&subscriptions::watching(&events.db, args(body)?).await?),
321 "inbox_watch" => reply(&subscriptions::watch(&events.db, args(body)?).await?),
322 "inbox_watchers" => reply(&subscriptions::watchers_count(&events.db, args(body)?).await?),
323 "inbox_watched" => reply(&subscriptions::watched(&events.db, args(body)?).await?),
324 "inbox_settings" => reply(&subscriptions::settings(&events.db, args(body)?).await?),
325 "inbox_update_settings" => reply(&subscriptions::update_settings(&events.db, args(body)?).await?),
326 _ => Response::error("Unknown method", 404),
327 }
328}
329
330/// Events from the bus. A batch that fails is retried whole, which is safe
331/// because the log ignores an event it has seen, and queues that already
332/// have it are skipped (`deliver`). After its retries it goes to the
333/// dead-letter queue (wrangler.jsonc).
334#[event(queue)]
335async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
336 let events = Events {
337 db: env.d1("DB")?,
338 env,
339 };
340 let delivered: Vec<Event> = batch
341 .messages()?
342 .into_iter()
343 .map(|message| message.into_body())
344 .collect();
345 events.deliver(&delivered).await?;
346 batch.ack_all();
347 Ok(())
348}
349
350/// Once a day, old audit entries are removed: first everything older than
351/// any workspace keeps (`AUDIT_MAX_DAYS`, 400 days), then each workspace's
352/// entries older than its own plan keeps, as billing says (7 days free, 90
353/// on the plan, or what staff set). Workspaces with nothing older than the
354/// shortest (`AUDIT_MIN_DAYS`, 7) are left alone.
355#[event(scheduled)]
356async fn scheduled(_event: worker::ScheduledEvent, env: Env, _ctx: worker::ScheduleContext) {
357 let days = |name: &str, default: u32| {
358 env.var(name)
359 .ok()
360 .and_then(|v| v.to_string().trim().parse::<u32>().ok())
361 .filter(|days| *days > 0)
362 .unwrap_or(default)
363 };
364 let max_days = days("AUDIT_MAX_DAYS", audit::DEFAULT_MAX_DAYS);
365 let min_days = days("AUDIT_MIN_DAYS", audit::DEFAULT_MIN_DAYS).min(max_days);
366 let Ok(db) = env.d1("DB") else { return };
367 let now = now_ms();
368 let cutoff = rfc3339(now.saturating_sub(FANOUT_KEEP_MS));
369 let purged = match db.prepare("DELETE FROM fanout_sent WHERE created_at < ?").bind(&[cutoff.into()]) {
370 Ok(statement) => statement.run().await.map(|_| ()),
371 Err(error) => Err(error),
372 };
373 if let Err(error) = purged {
374 worker::console_error!("could not remove old fan-out records: {error}");
375 }
376 match inbox::purge(&db, now).await {
377 Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"),
378 Ok(_) => {}
379 Err(error) => worker::console_error!("could not remove old inbox items: {error}"),
380 }
381 match audit::purge(&db, &audit::keep_from(now, max_days), 20).await {
382 Ok(removed) if removed > 0 => {
383 worker::console_log!("removed {removed} audit entries older than {max_days} days")
384 }
385 Ok(_) => {}
386 Err(error) => worker::console_error!("could not remove old audit entries: {error}"),
387 }
388 // Without billing nobody's plan is known, so nothing younger than the
389 // ceiling is removed.
390 let billing = match env.service("BILLING") {
391 Ok(billing) => billing,
392 Err(error) => {
393 worker::console_error!(
394 "audit entries kept past their plan's days: no billing: {error}"
395 );
396 return;
397 }
398 };
399 match audit::purge_by_plan(&db, &billing, now, min_days, max_days).await {
400 Ok(removed) if removed > 0 => {
401 worker::console_log!("removed {removed} audit entries older than their plan keeps")
402 }
403 Ok(_) => {}
404 Err(error) => worker::console_error!("audit entries kept past their plan's days: {error}"),
405 }
406}
407
408/// A value bound to a `?` in [`list_filter`]'s clause.
409#[derive(Debug, PartialEq)]
410enum Bind {
411 Text(String),
412 Number(f64),
413}
414
415/// The `WHERE` clause `list` reads with, and what it binds, in order.
416fn list_filter(a: &ListArgs) -> (String, Vec<Bind>) {
417 let mut conditions = Vec::new();
418 let mut values = Vec::new();
419 if let Some(repo_id) = &a.repo_id {
420 conditions.push("repo_id = ?".to_owned());
421 values.push(Bind::Text(repo_id.clone()));
422 }
423 if !a.types.is_empty() {
424 let marks = vec!["?"; a.types.len()].join(", ");
425 conditions.push(format!("type IN ({marks})"));
426 values.extend(a.types.iter().map(|kind| Bind::Text(kind.clone())));
427 } else {
428 // What CI and integrations report on commits goes to webhooks,
429 // and is read from each commit's checks; a timeline asked for
430 // everything would be little else on a busy repository.
431 let marks = vec!["?"; REPORTING.len()].join(", ");
432 conditions.push(format!("type NOT IN ({marks})"));
433 values.extend(REPORTING.iter().map(|kind| Bind::Text((*kind).to_owned())));
434 }
435 if let Some(actor) = &a.actor {
436 conditions.push("actor = ?".to_owned());
437 values.push(Bind::Text(actor.clone()));
438 }
439 if !a.numbers.is_empty() {
440 // An issue or pull request's own events name it as `number`; a
441 // comment, review or link on it names it as `issue`.
442 let marks = vec!["?"; a.numbers.len()].join(", ");
443 conditions.push(format!(
444 "(json_extract(data, '$.number') IN ({marks}) OR json_extract(data, '$.issue') IN ({marks}))"
445 ));
446 for _ in 0..2 {
447 values.extend(a.numbers.iter().map(|number| Bind::Number(f64::from(*number))));
448 }
449 }
450 if let Some(since) = &a.since {
451 conditions.push("time >= ?".to_owned());
452 values.push(Bind::Text(since.clone()));
453 }
454 if let Some(before) = &a.before {
455 conditions.push("id < ?".to_owned());
456 values.push(Bind::Text(before.clone()));
457 }
458 let filter = if conditions.is_empty() {
459 String::new()
460 } else {
461 format!("WHERE {}", conditions.join(" AND "))
462 };
463 (filter, values)
464}
465
466#[cfg(test)]
467mod tests {
468 use super::*;
469
470 #[test]
471 fn a_plain_list_leaves_out_what_is_reported_on_commits() {
472 let (filter, binds) = list_filter(&ListArgs { repo_id: Some("rep_1".into()), ..ListArgs::default() });
473 assert!(filter.starts_with("WHERE repo_id = ? AND type NOT IN ("));
474 assert_eq!(binds[0], Bind::Text("rep_1".into()));
475 assert_eq!(binds.len(), 1 + REPORTING.len());
476 }
477
478 #[test]
479 fn numbers_match_an_item_or_what_is_said_on_it_since_a_time() {
480 let (filter, binds) = list_filter(&ListArgs {
481 repo_id: Some("rep_1".into()),
482 types: vec!["issue.opened".into()],
483 numbers: vec![4, 9],
484 since: Some("2026-10-01T00:00:00Z".into()),
485 actor: Some("usr_g1t_agent".into()),
486 ..ListArgs::default()
487 });
488 assert_eq!(
489 filter,
490 "WHERE repo_id = ? AND type IN (?) AND actor = ? AND \
491 (json_extract(data, '$.number') IN (?, ?) OR json_extract(data, '$.issue') IN (?, ?)) AND time >= ?"
492 );
493 assert_eq!(
494 binds,
495 vec![
496 Bind::Text("rep_1".into()),
497 Bind::Text("issue.opened".into()),
498 Bind::Text("usr_g1t_agent".into()),
499 Bind::Number(4.0),
500 Bind::Number(9.0),
501 Bind::Number(4.0),
502 Bind::Number(9.0),
503 Bind::Text("2026-10-01T00:00:00Z".into()),
504 ]
505 );
506 }
507}