Skip to content
403 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 every subscriber's own
6//! queue, so a slow or failing subscriber holds up nobody else.
7//!
8//! It keeps two more things beside the log: the audit log (audit.rs) and
9//! each person's inbox (inbox.rs), written as events arrive, with who
10//! follows what (subscriptions.rs).
11//!
12//! Other services reach it over `POST /rpc/<method>`; see
13//! `g1t_contracts::events`, `audit` and `inbox` for the methods and their
14//! arguments.
15
16mod audit;
17mod inbox;
18mod subscriptions;
19
20use g1t_contracts::events::{Event, ListArgs, PublishArgs};
21use g1t_contracts::new_id;
22use g1t_contracts::time::rfc3339;
23use g1t_kit::{args, js, now_ms, reply, rpc_method};
24use serde::Deserialize;
25use worker::js_sys::{Array, Object};
26use worker::wasm_bindgen::{JsCast, JsValue};
27use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event};
28
29const DEFAULT_PAGE: u32 = 50;
30const MAX_PAGE: u32 = 200;
31/// Statuses and check runs reported on commits: delivered to webhooks, and
32/// left out of a timeline unless asked for by type.
33const REPORTING: [&str; 7] = [
34 "status.created",
35 "check_run.created",
36 "check_run.completed",
37 "check_run.rerequested",
38 "check_run.requested_action",
39 "check_suite.completed",
40 "check_suite.rerequested",
41];
42/// Every binding whose name starts with this is a queue that receives all
43/// events: one per subscribing service.
44const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_";
45
46#[derive(Deserialize)]
47struct EventRow {
48 id: String,
49 #[serde(rename = "type")]
50 kind: String,
51 source: String,
52 time: String,
53 repo_id: Option<String>,
54 actor: Option<String>,
55 /// JSON.
56 data: String,
57}
58
59impl From<EventRow> for Event {
60 fn from(row: EventRow) -> Self {
61 Event {
62 id: row.id,
63 kind: row.kind,
64 source: row.source,
65 time: row.time,
66 repo_id: row.repo_id,
67 actor: row.actor,
68 data: serde_json::from_str(&row.data).unwrap_or_default(),
69 }
70 }
71}
72
73fn optional(value: &Option<String>) -> JsValue {
74 value.as_deref().map_or(JsValue::NULL, JsValue::from)
75}
76
77/// Sends `events` to a queue binding, one message each.
78async fn send(queue: &JsValue, events: &[Event]) -> Result<()> {
79 let messages = Array::new();
80 for event in events {
81 let message = Object::new();
82 js::set(&message, "body", &js::to_js(event)?);
83 messages.push(&message);
84 }
85 js::call(queue, "sendBatch", &[messages.into()]).await?;
86 Ok(())
87}
88
89struct Events {
90 db: D1Database,
91 env: Env,
92}
93
94impl Events {
95 /// Assigns each event its id and time and puts it on the bus.
96 async fn publish(&self, a: PublishArgs) -> Result<()> {
97 if a.events.is_empty() {
98 return Ok(());
99 }
100 let now = now_ms();
101 let events: Vec<Event> = a
102 .events
103 .into_iter()
104 .map(|event| Event {
105 id: new_id("evt", now),
106 kind: event.kind,
107 source: event.source,
108 time: rfc3339(now),
109 repo_id: event.repo_id,
110 actor: event.actor,
111 data: event.data,
112 })
113 .collect();
114 send(&js::binding(&self.env, "BUS")?, &events).await
115 }
116
117 /// Newest first. Callers must have checked that the viewer may see the
118 /// repository asked about.
119 async fn list(&self, a: ListArgs) -> Result<Vec<Event>> {
120 let (filter, binds) = list_filter(&a);
121 let mut values: Vec<JsValue> = binds
122 .into_iter()
123 .map(|bind| match bind {
124 Bind::Text(text) => JsValue::from(text),
125 Bind::Number(number) => JsValue::from(number),
126 })
127 .collect();
128 values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into());
129 let rows = self
130 .db
131 .prepare(format!(
132 "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?"
133 ))
134 .bind(&values)?
135 .all()
136 .await?
137 .results::<EventRow>()?;
138 Ok(rows.into_iter().map(Event::from).collect())
139 }
140
141 /// Writes a batch from the bus to the log, hands it to every
142 /// subscriber, then tells the people it concerns (inbox.rs).
143 async fn deliver(&self, events: &[Event]) -> Result<()> {
144 let mut statements = Vec::with_capacity(events.len());
145 for event in events {
146 statements.push(
147 self.db
148 .prepare(
149 // Redelivered batches must not duplicate log rows.
150 "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data)
151 VALUES (?, ?, ?, ?, ?, ?, ?)",
152 )
153 .bind(&[
154 event.id.as_str().into(),
155 event.kind.as_str().into(),
156 event.source.as_str().into(),
157 event.time.as_str().into(),
158 optional(&event.repo_id),
159 optional(&event.actor),
160 serde_json::to_string(&event.data)?.into(),
161 ])?,
162 );
163 }
164 self.db.batch(statements).await?;
165 audit::follow_renames(&self.db, events).await?;
166
167 let bindings: &JsValue = self.env.as_ref();
168 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
169 let Some(name) = name
170 .as_string()
171 .filter(|name| name.starts_with(SUBSCRIBER_PREFIX))
172 else {
173 continue;
174 };
175 send(&js::binding(&self.env, &name)?, events).await?;
176 }
177 inbox::follow(&self.db, events).await?;
178 let (work, repos, identity) = (
179 self.env.service("WORK")?,
180 self.env.service("REPOS")?,
181 self.env.service("IDENTITY")?,
182 );
183 let sources = inbox::Sources {
184 work: &work,
185 repos: &repos,
186 identity: &identity,
187 };
188 inbox::deliver(&self.db, &sources, events).await;
189 Ok(())
190 }
191}
192
193#[event(fetch)]
194async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
195 let Some(method) = rpc_method(&request) else {
196 return Response::error("Not found", 404);
197 };
198 let body: serde_json::Value = request.json().await?;
199 let events = Events {
200 db: env.d1("DB")?,
201 env,
202 };
203 match method.as_str() {
204 "publish" => reply(&events.publish(args(body)?).await?),
205 "list" => reply(&events.list(args(body)?).await?),
206 "audit_record" => reply(&audit::record(&events.db, args(body)?).await?),
207 "audit_list" => reply(&audit::list(&events.db, args(body)?).await?),
208 "inbox_list" => {
209 let repos = events.env.service("REPOS")?;
210 reply(&inbox::list(&events.db, &repos, args(body)?).await?)
211 }
212 "inbox_counts" => reply(&inbox::counts(&events.db, args(body)?).await?),
213 "inbox_mark" => reply(&inbox::mark(&events.db, args(body)?).await?),
214 "inbox_thread" => {
215 let (repos, work) = (events.env.service("REPOS")?, events.env.service("WORK")?);
216 reply(&inbox::thread(&events.db, &repos, &work, args(body)?).await?)
217 }
218 "inbox_subscription" => {
219 let work = events.env.service("WORK")?;
220 reply(&subscriptions::subscription(&events.db, &work, args(body)?).await?)
221 }
222 "inbox_subscribe" => {
223 let work = events.env.service("WORK")?;
224 reply(&subscriptions::subscribe(&events.db, &work, args(body)?).await?)
225 }
226 "inbox_watching" => reply(&subscriptions::watching(&events.db, args(body)?).await?),
227 "inbox_watch" => reply(&subscriptions::watch(&events.db, args(body)?).await?),
228 "inbox_watchers" => reply(&subscriptions::watchers_count(&events.db, args(body)?).await?),
229 "inbox_watched" => reply(&subscriptions::watched(&events.db, args(body)?).await?),
230 "inbox_settings" => reply(&subscriptions::settings(&events.db, args(body)?).await?),
231 "inbox_update_settings" => reply(&subscriptions::update_settings(&events.db, args(body)?).await?),
232 _ => Response::error("Unknown method", 404),
233 }
234}
235
236/// Events from the bus. A batch that fails is retried whole, which is safe
237/// because both the log and subscribers ignore an event they have seen.
238#[event(queue)]
239async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
240 let events = Events {
241 db: env.d1("DB")?,
242 env,
243 };
244 let delivered: Vec<Event> = batch
245 .messages()?
246 .into_iter()
247 .map(|message| message.into_body())
248 .collect();
249 events.deliver(&delivered).await?;
250 batch.ack_all();
251 Ok(())
252}
253
254/// Once a day, old audit entries are removed: first everything older than
255/// any workspace keeps (`AUDIT_MAX_DAYS`, 400 days), then each workspace's
256/// entries older than its own plan keeps, as billing says (7 days free, 90
257/// on the plan, or what staff set). Workspaces with nothing older than the
258/// shortest (`AUDIT_MIN_DAYS`, 7) are left alone.
259#[event(scheduled)]
260async fn scheduled(_event: worker::ScheduledEvent, env: Env, _ctx: worker::ScheduleContext) {
261 let days = |name: &str, default: u32| {
262 env.var(name)
263 .ok()
264 .and_then(|v| v.to_string().trim().parse::<u32>().ok())
265 .filter(|days| *days > 0)
266 .unwrap_or(default)
267 };
268 let max_days = days("AUDIT_MAX_DAYS", audit::DEFAULT_MAX_DAYS);
269 let min_days = days("AUDIT_MIN_DAYS", audit::DEFAULT_MIN_DAYS).min(max_days);
270 let Ok(db) = env.d1("DB") else { return };
271 let now = now_ms();
272 match inbox::purge(&db, now).await {
273 Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"),
274 Ok(_) => {}
275 Err(error) => worker::console_error!("could not remove old inbox items: {error}"),
276 }
277 match audit::purge(&db, &audit::keep_from(now, max_days), 20).await {
278 Ok(removed) if removed > 0 => {
279 worker::console_log!("removed {removed} audit entries older than {max_days} days")
280 }
281 Ok(_) => {}
282 Err(error) => worker::console_error!("could not remove old audit entries: {error}"),
283 }
284 // Without billing nobody's plan is known, so nothing younger than the
285 // ceiling is removed.
286 let billing = match env.service("BILLING") {
287 Ok(billing) => billing,
288 Err(error) => {
289 worker::console_error!(
290 "audit entries kept past their plan's days: no billing: {error}"
291 );
292 return;
293 }
294 };
295 match audit::purge_by_plan(&db, &billing, now, min_days, max_days).await {
296 Ok(removed) if removed > 0 => {
297 worker::console_log!("removed {removed} audit entries older than their plan keeps")
298 }
299 Ok(_) => {}
300 Err(error) => worker::console_error!("audit entries kept past their plan's days: {error}"),
301 }
302}
303
304/// A value bound to a `?` in [`list_filter`]'s clause.
305#[derive(Debug, PartialEq)]
306enum Bind {
307 Text(String),
308 Number(f64),
309}
310
311/// The `WHERE` clause `list` reads with, and what it binds, in order.
312fn list_filter(a: &ListArgs) -> (String, Vec<Bind>) {
313 let mut conditions = Vec::new();
314 let mut values = Vec::new();
315 if let Some(repo_id) = &a.repo_id {
316 conditions.push("repo_id = ?".to_owned());
317 values.push(Bind::Text(repo_id.clone()));
318 }
319 if !a.types.is_empty() {
320 let marks = vec!["?"; a.types.len()].join(", ");
321 conditions.push(format!("type IN ({marks})"));
322 values.extend(a.types.iter().map(|kind| Bind::Text(kind.clone())));
323 } else {
324 // What CI and integrations report on commits goes to webhooks,
325 // and is read from each commit's checks; a timeline asked for
326 // everything would be little else on a busy repository.
327 let marks = vec!["?"; REPORTING.len()].join(", ");
328 conditions.push(format!("type NOT IN ({marks})"));
329 values.extend(REPORTING.iter().map(|kind| Bind::Text((*kind).to_owned())));
330 }
331 if let Some(actor) = &a.actor {
332 conditions.push("actor = ?".to_owned());
333 values.push(Bind::Text(actor.clone()));
334 }
335 if !a.numbers.is_empty() {
336 // An issue or pull request's own events name it as `number`; a
337 // comment, review or link on it names it as `issue`.
338 let marks = vec!["?"; a.numbers.len()].join(", ");
339 conditions.push(format!(
340 "(json_extract(data, '$.number') IN ({marks}) OR json_extract(data, '$.issue') IN ({marks}))"
341 ));
342 for _ in 0..2 {
343 values.extend(a.numbers.iter().map(|number| Bind::Number(f64::from(*number))));
344 }
345 }
346 if let Some(since) = &a.since {
347 conditions.push("time >= ?".to_owned());
348 values.push(Bind::Text(since.clone()));
349 }
350 if let Some(before) = &a.before {
351 conditions.push("id < ?".to_owned());
352 values.push(Bind::Text(before.clone()));
353 }
354 let filter = if conditions.is_empty() {
355 String::new()
356 } else {
357 format!("WHERE {}", conditions.join(" AND "))
358 };
359 (filter, values)
360}
361
362#[cfg(test)]
363mod tests {
364 use super::*;
365
366 #[test]
367 fn a_plain_list_leaves_out_what_is_reported_on_commits() {
368 let (filter, binds) = list_filter(&ListArgs { repo_id: Some("rep_1".into()), ..ListArgs::default() });
369 assert!(filter.starts_with("WHERE repo_id = ? AND type NOT IN ("));
370 assert_eq!(binds[0], Bind::Text("rep_1".into()));
371 assert_eq!(binds.len(), 1 + REPORTING.len());
372 }
373
374 #[test]
375 fn numbers_match_an_item_or_what_is_said_on_it_since_a_time() {
376 let (filter, binds) = list_filter(&ListArgs {
377 repo_id: Some("rep_1".into()),
378 types: vec!["issue.opened".into()],
379 numbers: vec![4, 9],
380 since: Some("2026-10-01T00:00:00Z".into()),
381 actor: Some("usr_g1t_agent".into()),
382 ..ListArgs::default()
383 });
384 assert_eq!(
385 filter,
386 "WHERE repo_id = ? AND type IN (?) AND actor = ? AND \
387 (json_extract(data, '$.number') IN (?, ?) OR json_extract(data, '$.issue') IN (?, ?)) AND time >= ?"
388 );
389 assert_eq!(
390 binds,
391 vec![
392 Bind::Text("rep_1".into()),
393 Bind::Text("issue.opened".into()),
394 Bind::Text("usr_g1t_agent".into()),
395 Bind::Number(4.0),
396 Bind::Number(9.0),
397 Bind::Number(4.0),
398 Bind::Number(9.0),
399 Bind::Text("2026-10-01T00:00:00Z".into()),
400 ]
401 );
402 }
403}