Skip to content

g1t/services/events/src/lib.rs

303 lines11,482 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/// Every binding whose name starts with this is a queue that receives all
32/// events: one per subscribing service.
33const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_";
34
35#[derive(Deserialize)]
36struct EventRow {
37 id: String,
38 #[serde(rename = "type")]
39 kind: String,
40 source: String,
41 time: String,
42 repo_id: Option<String>,
43 actor: Option<String>,
44 /// JSON.
45 data: String,
46}
47
48impl From<EventRow> for Event {
49 fn from(row: EventRow) -> Self {
50 Event {
51 id: row.id,
52 kind: row.kind,
53 source: row.source,
54 time: row.time,
55 repo_id: row.repo_id,
56 actor: row.actor,
57 data: serde_json::from_str(&row.data).unwrap_or_default(),
58 }
59 }
60}
61
62fn optional(value: &Option<String>) -> JsValue {
63 value.as_deref().map_or(JsValue::NULL, JsValue::from)
64}
65
66/// Sends `events` to a queue binding, one message each.
67async fn send(queue: &JsValue, events: &[Event]) -> Result<()> {
68 let messages = Array::new();
69 for event in events {
70 let message = Object::new();
71 js::set(&message, "body", &js::to_js(event)?);
72 messages.push(&message);
73 }
74 js::call(queue, "sendBatch", &[messages.into()]).await?;
75 Ok(())
76}
77
78struct Events {
79 db: D1Database,
80 env: Env,
81}
82
83impl Events {
84 /// Assigns each event its id and time and puts it on the bus.
85 async fn publish(&self, a: PublishArgs) -> Result<()> {
86 if a.events.is_empty() {
87 return Ok(());
88 }
89 let now = now_ms();
90 let events: Vec<Event> = a
91 .events
92 .into_iter()
93 .map(|event| Event {
94 id: new_id("evt", now),
95 kind: event.kind,
96 source: event.source,
97 time: rfc3339(now),
98 repo_id: event.repo_id,
99 actor: event.actor,
100 data: event.data,
101 })
102 .collect();
103 send(&js::binding(&self.env, "BUS")?, &events).await
104 }
105
106 /// Newest first. Callers must have checked that the viewer may see the
107 /// repository asked about.
108 async fn list(&self, a: ListArgs) -> Result<Vec<Event>> {
109 let mut conditions = Vec::new();
110 let mut values: Vec<JsValue> = Vec::new();
111 if let Some(repo_id) = &a.repo_id {
112 conditions.push("repo_id = ?".to_owned());
113 values.push(repo_id.as_str().into());
114 }
115 if !a.types.is_empty() {
116 let marks = vec!["?"; a.types.len()].join(", ");
117 conditions.push(format!("type IN ({marks})"));
118 values.extend(a.types.iter().map(|kind| JsValue::from(kind.as_str())));
119 }
120 if let Some(before) = &a.before {
121 conditions.push("id < ?".to_owned());
122 values.push(before.as_str().into());
123 }
124 let filter = if conditions.is_empty() {
125 String::new()
126 } else {
127 format!("WHERE {}", conditions.join(" AND "))
128 };
129 values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into());
130 let rows = self
131 .db
132 .prepare(format!(
133 "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?"
134 ))
135 .bind(&values)?
136 .all()
137 .await?
138 .results::<EventRow>()?;
139 Ok(rows.into_iter().map(Event::from).collect())
140 }
141
142 /// Writes a batch from the bus to the log, hands it to every
143 /// subscriber, then tells the people it concerns (inbox.rs).
144 async fn deliver(&self, events: &[Event]) -> Result<()> {
145 let mut statements = Vec::with_capacity(events.len());
146 for event in events {
147 statements.push(
148 self.db
149 .prepare(
150 // Redelivered batches must not duplicate log rows.
151 "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data)
152 VALUES (?, ?, ?, ?, ?, ?, ?)",
153 )
154 .bind(&[
155 event.id.as_str().into(),
156 event.kind.as_str().into(),
157 event.source.as_str().into(),
158 event.time.as_str().into(),
159 optional(&event.repo_id),
160 optional(&event.actor),
161 serde_json::to_string(&event.data)?.into(),
162 ])?,
163 );
164 }
165 self.db.batch(statements).await?;
166 audit::follow_renames(&self.db, events).await?;
167
168 let bindings: &JsValue = self.env.as_ref();
169 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
170 let Some(name) = name
171 .as_string()
172 .filter(|name| name.starts_with(SUBSCRIBER_PREFIX))
173 else {
174 continue;
175 };
176 send(&js::binding(&self.env, &name)?, events).await?;
177 }
178 inbox::follow(&self.db, events).await?;
179 let (work, repos, identity) = (
180 self.env.service("WORK")?,
181 self.env.service("REPOS")?,
182 self.env.service("IDENTITY")?,
183 );
184 let sources = inbox::Sources {
185 work: &work,
186 repos: &repos,
187 identity: &identity,
188 };
189 inbox::deliver(&self.db, &sources, events).await;
190 Ok(())
191 }
192}
193
194#[event(fetch)]
195async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
196 let Some(method) = rpc_method(&request) else {
197 return Response::error("Not found", 404);
198 };
199 let body: serde_json::Value = request.json().await?;
200 let events = Events {
201 db: env.d1("DB")?,
202 env,
203 };
204 match method.as_str() {
205 "publish" => reply(&events.publish(args(body)?).await?),
206 "list" => reply(&events.list(args(body)?).await?),
207 "audit_record" => reply(&audit::record(&events.db, args(body)?).await?),
208 "audit_list" => reply(&audit::list(&events.db, args(body)?).await?),
209 "inbox_list" => {
210 let repos = events.env.service("REPOS")?;
211 reply(&inbox::list(&events.db, &repos, args(body)?).await?)
212 }
213 "inbox_counts" => reply(&inbox::counts(&events.db, args(body)?).await?),
214 "inbox_mark" => reply(&inbox::mark(&events.db, args(body)?).await?),
215 "inbox_thread" => {
216 let (repos, work) = (events.env.service("REPOS")?, events.env.service("WORK")?);
217 reply(&inbox::thread(&events.db, &repos, &work, args(body)?).await?)
218 }
219 "inbox_subscription" => {
220 let work = events.env.service("WORK")?;
221 reply(&subscriptions::subscription(&events.db, &work, args(body)?).await?)
222 }
223 "inbox_subscribe" => {
224 let work = events.env.service("WORK")?;
225 reply(&subscriptions::subscribe(&events.db, &work, args(body)?).await?)
226 }
227 "inbox_watching" => reply(&subscriptions::watching(&events.db, args(body)?).await?),
228 "inbox_watch" => reply(&subscriptions::watch(&events.db, args(body)?).await?),
229 "inbox_watchers" => reply(&subscriptions::watchers_count(&events.db, args(body)?).await?),
230 "inbox_watched" => reply(&subscriptions::watched(&events.db, args(body)?).await?),
231 "inbox_settings" => reply(&subscriptions::settings(&events.db, args(body)?).await?),
232 "inbox_update_settings" => reply(&subscriptions::update_settings(&events.db, args(body)?).await?),
233 _ => Response::error("Unknown method", 404),
234 }
235}
236
237/// Events from the bus. A batch that fails is retried whole, which is safe
238/// because both the log and subscribers ignore an event they have seen.
239#[event(queue)]
240async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
241 let events = Events {
242 db: env.d1("DB")?,
243 env,
244 };
245 let delivered: Vec<Event> = batch
246 .messages()?
247 .into_iter()
248 .map(|message| message.into_body())
249 .collect();
250 events.deliver(&delivered).await?;
251 batch.ack_all();
252 Ok(())
253}
254
255/// Once a day, old audit entries are removed: first everything older than
256/// any workspace keeps (`AUDIT_MAX_DAYS`, 400 days), then each workspace's
257/// entries older than its own plan keeps, as billing says (7 days free, 90
258/// on the plan, or what staff set). Workspaces with nothing older than the
259/// shortest (`AUDIT_MIN_DAYS`, 7) are left alone.
260#[event(scheduled)]
261async fn scheduled(_event: worker::ScheduledEvent, env: Env, _ctx: worker::ScheduleContext) {
262 let days = |name: &str, default: u32| {
263 env.var(name)
264 .ok()
265 .and_then(|v| v.to_string().trim().parse::<u32>().ok())
266 .filter(|days| *days > 0)
267 .unwrap_or(default)
268 };
269 let max_days = days("AUDIT_MAX_DAYS", audit::DEFAULT_MAX_DAYS);
270 let min_days = days("AUDIT_MIN_DAYS", audit::DEFAULT_MIN_DAYS).min(max_days);
271 let Ok(db) = env.d1("DB") else { return };
272 let now = now_ms();
273 match inbox::purge(&db, now).await {
274 Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"),
275 Ok(_) => {}
276 Err(error) => worker::console_error!("could not remove old inbox items: {error}"),
277 }
278 match audit::purge(&db, &audit::keep_from(now, max_days), 20).await {
279 Ok(removed) if removed > 0 => {
280 worker::console_log!("removed {removed} audit entries older than {max_days} days")
281 }
282 Ok(_) => {}
283 Err(error) => worker::console_error!("could not remove old audit entries: {error}"),
284 }
285 // Without billing nobody's plan is known, so nothing younger than the
286 // ceiling is removed.
287 let billing = match env.service("BILLING") {
288 Ok(billing) => billing,
289 Err(error) => {
290 worker::console_error!(
291 "audit entries kept past their plan's days: no billing: {error}"
292 );
293 return;
294 }
295 };
296 match audit::purge_by_plan(&db, &billing, now, min_days, max_days).await {
297 Ok(removed) if removed > 0 => {
298 worker::console_log!("removed {removed} audit entries older than their plan keeps")
299 }
300 Ok(_) => {}
301 Err(error) => worker::console_error!("audit entries kept past their plan's days: {error}"),
302 }
303}