Skip to content

g1t/services/events/src/lib.rs

321 lines12,283 bytesCodeBlame
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 mut conditions = Vec::new();
121 let mut values: Vec<JsValue> = Vec::new();
122 if let Some(repo_id) = &a.repo_id {
123 conditions.push("repo_id = ?".to_owned());
124 values.push(repo_id.as_str().into());
125 }
126 if !a.types.is_empty() {
127 let marks = vec!["?"; a.types.len()].join(", ");
128 conditions.push(format!("type IN ({marks})"));
129 values.extend(a.types.iter().map(|kind| JsValue::from(kind.as_str())));
130 } else {
131 // What CI and integrations report on commits goes to webhooks,
132 // and is read from each commit's checks; a timeline asked for
133 // everything would be little else on a busy repository.
134 let marks = vec!["?"; REPORTING.len()].join(", ");
135 conditions.push(format!("type NOT IN ({marks})"));
136 values.extend(REPORTING.iter().map(|kind| JsValue::from(*kind)));
137 }
138 if let Some(before) = &a.before {
139 conditions.push("id < ?".to_owned());
140 values.push(before.as_str().into());
141 }
142 let filter = if conditions.is_empty() {
143 String::new()
144 } else {
145 format!("WHERE {}", conditions.join(" AND "))
146 };
147 values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into());
148 let rows = self
149 .db
150 .prepare(format!(
151 "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?"
152 ))
153 .bind(&values)?
154 .all()
155 .await?
156 .results::<EventRow>()?;
157 Ok(rows.into_iter().map(Event::from).collect())
158 }
159
160 /// Writes a batch from the bus to the log, hands it to every
161 /// subscriber, then tells the people it concerns (inbox.rs).
162 async fn deliver(&self, events: &[Event]) -> Result<()> {
163 let mut statements = Vec::with_capacity(events.len());
164 for event in events {
165 statements.push(
166 self.db
167 .prepare(
168 // Redelivered batches must not duplicate log rows.
169 "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data)
170 VALUES (?, ?, ?, ?, ?, ?, ?)",
171 )
172 .bind(&[
173 event.id.as_str().into(),
174 event.kind.as_str().into(),
175 event.source.as_str().into(),
176 event.time.as_str().into(),
177 optional(&event.repo_id),
178 optional(&event.actor),
179 serde_json::to_string(&event.data)?.into(),
180 ])?,
181 );
182 }
183 self.db.batch(statements).await?;
184 audit::follow_renames(&self.db, events).await?;
185
186 let bindings: &JsValue = self.env.as_ref();
187 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
188 let Some(name) = name
189 .as_string()
190 .filter(|name| name.starts_with(SUBSCRIBER_PREFIX))
191 else {
192 continue;
193 };
194 send(&js::binding(&self.env, &name)?, events).await?;
195 }
196 inbox::follow(&self.db, events).await?;
197 let (work, repos, identity) = (
198 self.env.service("WORK")?,
199 self.env.service("REPOS")?,
200 self.env.service("IDENTITY")?,
201 );
202 let sources = inbox::Sources {
203 work: &work,
204 repos: &repos,
205 identity: &identity,
206 };
207 inbox::deliver(&self.db, &sources, events).await;
208 Ok(())
209 }
210}
211
212#[event(fetch)]
213async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
214 let Some(method) = rpc_method(&request) else {
215 return Response::error("Not found", 404);
216 };
217 let body: serde_json::Value = request.json().await?;
218 let events = Events {
219 db: env.d1("DB")?,
220 env,
221 };
222 match method.as_str() {
223 "publish" => reply(&events.publish(args(body)?).await?),
224 "list" => reply(&events.list(args(body)?).await?),
225 "audit_record" => reply(&audit::record(&events.db, args(body)?).await?),
226 "audit_list" => reply(&audit::list(&events.db, args(body)?).await?),
227 "inbox_list" => {
228 let repos = events.env.service("REPOS")?;
229 reply(&inbox::list(&events.db, &repos, args(body)?).await?)
230 }
231 "inbox_counts" => reply(&inbox::counts(&events.db, args(body)?).await?),
232 "inbox_mark" => reply(&inbox::mark(&events.db, args(body)?).await?),
233 "inbox_thread" => {
234 let (repos, work) = (events.env.service("REPOS")?, events.env.service("WORK")?);
235 reply(&inbox::thread(&events.db, &repos, &work, args(body)?).await?)
236 }
237 "inbox_subscription" => {
238 let work = events.env.service("WORK")?;
239 reply(&subscriptions::subscription(&events.db, &work, args(body)?).await?)
240 }
241 "inbox_subscribe" => {
242 let work = events.env.service("WORK")?;
243 reply(&subscriptions::subscribe(&events.db, &work, args(body)?).await?)
244 }
245 "inbox_watching" => reply(&subscriptions::watching(&events.db, args(body)?).await?),
246 "inbox_watch" => reply(&subscriptions::watch(&events.db, args(body)?).await?),
247 "inbox_watchers" => reply(&subscriptions::watchers_count(&events.db, args(body)?).await?),
248 "inbox_watched" => reply(&subscriptions::watched(&events.db, args(body)?).await?),
249 "inbox_settings" => reply(&subscriptions::settings(&events.db, args(body)?).await?),
250 "inbox_update_settings" => reply(&subscriptions::update_settings(&events.db, args(body)?).await?),
251 _ => Response::error("Unknown method", 404),
252 }
253}
254
255/// Events from the bus. A batch that fails is retried whole, which is safe
256/// because both the log and subscribers ignore an event they have seen.
257#[event(queue)]
258async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
259 let events = Events {
260 db: env.d1("DB")?,
261 env,
262 };
263 let delivered: Vec<Event> = batch
264 .messages()?
265 .into_iter()
266 .map(|message| message.into_body())
267 .collect();
268 events.deliver(&delivered).await?;
269 batch.ack_all();
270 Ok(())
271}
272
273/// Once a day, old audit entries are removed: first everything older than
274/// any workspace keeps (`AUDIT_MAX_DAYS`, 400 days), then each workspace's
275/// entries older than its own plan keeps, as billing says (7 days free, 90
276/// on the plan, or what staff set). Workspaces with nothing older than the
277/// shortest (`AUDIT_MIN_DAYS`, 7) are left alone.
278#[event(scheduled)]
279async fn scheduled(_event: worker::ScheduledEvent, env: Env, _ctx: worker::ScheduleContext) {
280 let days = |name: &str, default: u32| {
281 env.var(name)
282 .ok()
283 .and_then(|v| v.to_string().trim().parse::<u32>().ok())
284 .filter(|days| *days > 0)
285 .unwrap_or(default)
286 };
287 let max_days = days("AUDIT_MAX_DAYS", audit::DEFAULT_MAX_DAYS);
288 let min_days = days("AUDIT_MIN_DAYS", audit::DEFAULT_MIN_DAYS).min(max_days);
289 let Ok(db) = env.d1("DB") else { return };
290 let now = now_ms();
291 match inbox::purge(&db, now).await {
292 Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"),
293 Ok(_) => {}
294 Err(error) => worker::console_error!("could not remove old inbox items: {error}"),
295 }
296 match audit::purge(&db, &audit::keep_from(now, max_days), 20).await {
297 Ok(removed) if removed > 0 => {
298 worker::console_log!("removed {removed} audit entries older than {max_days} days")
299 }
300 Ok(_) => {}
301 Err(error) => worker::console_error!("could not remove old audit entries: {error}"),
302 }
303 // Without billing nobody's plan is known, so nothing younger than the
304 // ceiling is removed.
305 let billing = match env.service("BILLING") {
306 Ok(billing) => billing,
307 Err(error) => {
308 worker::console_error!(
309 "audit entries kept past their plan's days: no billing: {error}"
310 );
311 return;
312 }
313 };
314 match audit::purge_by_plan(&db, &billing, now, min_days, max_days).await {
315 Ok(removed) if removed > 0 => {
316 worker::console_log!("removed {removed} audit entries older than their plan keeps")
317 }
318 Ok(_) => {}
319 Err(error) => worker::console_error!("audit entries kept past their plan's days: {error}"),
320 }
321}