Skip to content

g1t/services/events/src/lib.rs

283 lines10,275 bytesCodeBlame

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.

Events service in Rust, with RFC 3339 times and accurate push events1//! 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//!
Inbox: the events service tells people what needs them as events arrive8//! 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.
10//!
Events service in Rust, with RFC 3339 times and accurate push events11//! Other services reach it over `POST /rpc/<method>`; see
Inbox: the events service tells people what needs them as events arrive12//! `g1t_contracts::events`, `audit` and `inbox` for the methods and their
13//! arguments.
Events service in Rust, with RFC 3339 times and accurate push events14
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API15mod audit;
Inbox: the events service tells people what needs them as events arrive16mod inbox;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API17
Events service in Rust, with RFC 3339 times and accurate push events18use g1t_contracts::events::{Event, ListArgs, PublishArgs};
19use g1t_contracts::new_id;
20use g1t_contracts::time::rfc3339;
21use g1t_kit::{args, js, now_ms, reply, rpc_method};
22use serde::Deserialize;
23use worker::js_sys::{Array, Object};
24use worker::wasm_bindgen::{JsCast, JsValue};
25use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event};
26
27const DEFAULT_PAGE: u32 = 50;
28const MAX_PAGE: u32 = 200;
29/// Every binding whose name starts with this is a queue that receives all
30/// events: one per subscribing service.
31const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_";
32
33#[derive(Deserialize)]
34struct EventRow {
35 id: String,
36 #[serde(rename = "type")]
37 kind: String,
38 source: String,
39 time: String,
40 repo_id: Option<String>,
41 actor: Option<String>,
42 /// JSON.
43 data: String,
44}
45
46impl From<EventRow> for Event {
47 fn from(row: EventRow) -> Self {
48 Event {
49 id: row.id,
50 kind: row.kind,
51 source: row.source,
52 time: row.time,
53 repo_id: row.repo_id,
54 actor: row.actor,
55 data: serde_json::from_str(&row.data).unwrap_or_default(),
56 }
57 }
58}
59
60fn optional(value: &Option<String>) -> JsValue {
61 value.as_deref().map_or(JsValue::NULL, JsValue::from)
62}
63
64/// Sends `events` to a queue binding, one message each.
65async fn send(queue: &JsValue, events: &[Event]) -> Result<()> {
66 let messages = Array::new();
67 for event in events {
68 let message = Object::new();
69 js::set(&message, "body", &js::to_js(event)?);
70 messages.push(&message);
71 }
72 js::call(queue, "sendBatch", &[messages.into()]).await?;
73 Ok(())
74}
75
76struct Events {
77 db: D1Database,
78 env: Env,
79}
80
81impl Events {
82 /// Assigns each event its id and time and puts it on the bus.
83 async fn publish(&self, a: PublishArgs) -> Result<()> {
84 if a.events.is_empty() {
85 return Ok(());
86 }
87 let now = now_ms();
88 let events: Vec<Event> = a
89 .events
90 .into_iter()
91 .map(|event| Event {
92 id: new_id("evt", now),
93 kind: event.kind,
94 source: event.source,
95 time: rfc3339(now),
96 repo_id: event.repo_id,
97 actor: event.actor,
98 data: event.data,
99 })
100 .collect();
101 send(&js::binding(&self.env, "BUS")?, &events).await
102 }
103
104 /// Newest first. Callers must have checked that the viewer may see the
105 /// repository asked about.
106 async fn list(&self, a: ListArgs) -> Result<Vec<Event>> {
107 let mut conditions = Vec::new();
108 let mut values: Vec<JsValue> = Vec::new();
109 if let Some(repo_id) = &a.repo_id {
110 conditions.push("repo_id = ?".to_owned());
111 values.push(repo_id.as_str().into());
112 }
113 if !a.types.is_empty() {
114 let marks = vec!["?"; a.types.len()].join(", ");
115 conditions.push(format!("type IN ({marks})"));
116 values.extend(a.types.iter().map(|kind| JsValue::from(kind.as_str())));
117 }
118 if let Some(before) = &a.before {
119 conditions.push("id < ?".to_owned());
120 values.push(before.as_str().into());
121 }
122 let filter = if conditions.is_empty() {
123 String::new()
124 } else {
125 format!("WHERE {}", conditions.join(" AND "))
126 };
127 values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into());
128 let rows = self
129 .db
130 .prepare(format!(
131 "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?"
132 ))
133 .bind(&values)?
134 .all()
135 .await?
136 .results::<EventRow>()?;
137 Ok(rows.into_iter().map(Event::from).collect())
138 }
139
Inbox: the events service tells people what needs them as events arrive140 /// Writes a batch from the bus to the log, hands it to every
141 /// subscriber, then tells the people it concerns (inbox.rs).
Events service in Rust, with RFC 3339 times and accurate push events142 async fn deliver(&self, events: &[Event]) -> Result<()> {
143 let mut statements = Vec::with_capacity(events.len());
144 for event in events {
145 statements.push(
146 self.db
147 .prepare(
148 // Redelivered batches must not duplicate log rows.
149 "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data)
150 VALUES (?, ?, ?, ?, ?, ?, ?)",
151 )
152 .bind(&[
153 event.id.as_str().into(),
154 event.kind.as_str().into(),
155 event.source.as_str().into(),
156 event.time.as_str().into(),
157 optional(&event.repo_id),
158 optional(&event.actor),
159 serde_json::to_string(&event.data)?.into(),
160 ])?,
161 );
162 }
163 self.db.batch(statements).await?;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API164 audit::follow_renames(&self.db, events).await?;
Events service in Rust, with RFC 3339 times and accurate push events165
166 let bindings: &JsValue = self.env.as_ref();
167 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
168 let Some(name) = name
169 .as_string()
170 .filter(|name| name.starts_with(SUBSCRIBER_PREFIX))
171 else {
172 continue;
173 };
174 send(&js::binding(&self.env, &name)?, events).await?;
175 }
Inbox: the events service tells people what needs them as events arrive176 inbox::follow(&self.db, events).await?;
177 let (work, repos, identity) = (
178 self.env.service("WORK")?,
179 self.env.service("REPOS")?,
180 self.env.service("IDENTITY")?,
181 );
182 let sources = inbox::Sources {
183 work: &work,
184 repos: &repos,
185 identity: &identity,
186 };
187 inbox::deliver(&self.db, &sources, events).await;
Events service in Rust, with RFC 3339 times and accurate push events188 Ok(())
189 }
190}
191
192#[event(fetch)]
193async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
194 let Some(method) = rpc_method(&request) else {
195 return Response::error("Not found", 404);
196 };
197 let body: serde_json::Value = request.json().await?;
198 let events = Events {
199 db: env.d1("DB")?,
200 env,
201 };
202 match method.as_str() {
203 "publish" => reply(&events.publish(args(body)?).await?),
204 "list" => reply(&events.list(args(body)?).await?),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API205 "audit_record" => reply(&audit::record(&events.db, args(body)?).await?),
206 "audit_list" => reply(&audit::list(&events.db, args(body)?).await?),
Inbox: the events service tells people what needs them as events arrive207 "inbox_list" => {
208 let repos = events.env.service("REPOS")?;
209 reply(&inbox::list(&events.db, &repos, args(body)?).await?)
210 }
211 "inbox_counts" => reply(&inbox::counts(&events.db, args(body)?).await?),
212 "inbox_mark" => reply(&inbox::mark(&events.db, args(body)?).await?),
Events service in Rust, with RFC 3339 times and accurate push events213 _ => Response::error("Unknown method", 404),
214 }
215}
216
217/// Events from the bus. A batch that fails is retried whole, which is safe
218/// because both the log and subscribers ignore an event they have seen.
219#[event(queue)]
220async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
221 let events = Events {
222 db: env.d1("DB")?,
223 env,
224 };
225 let delivered: Vec<Event> = batch
226 .messages()?
227 .into_iter()
228 .map(|message| message.into_body())
229 .collect();
230 events.deliver(&delivered).await?;
231 batch.ack_all();
232 Ok(())
233}
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put234
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo235/// Once a day, old audit entries are removed: first everything older than
236/// any workspace keeps (`AUDIT_MAX_DAYS`, 400 days), then each workspace's
237/// entries older than its own plan keeps, as billing says (7 days free, 90
238/// on the plan, or what staff set). Workspaces with nothing older than the
239/// shortest (`AUDIT_MIN_DAYS`, 7) are left alone.
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put240#[event(scheduled)]
241async fn scheduled(_event: worker::ScheduledEvent, env: Env, _ctx: worker::ScheduleContext) {
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo242 let days = |name: &str, default: u32| {
243 env.var(name)
244 .ok()
245 .and_then(|v| v.to_string().trim().parse::<u32>().ok())
246 .filter(|days| *days > 0)
247 .unwrap_or(default)
248 };
249 let max_days = days("AUDIT_MAX_DAYS", audit::DEFAULT_MAX_DAYS);
250 let min_days = days("AUDIT_MIN_DAYS", audit::DEFAULT_MIN_DAYS).min(max_days);
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put251 let Ok(db) = env.d1("DB") else { return };
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo252 let now = now_ms();
Inbox: the events service tells people what needs them as events arrive253 match inbox::purge(&db, now).await {
254 Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"),
255 Ok(_) => {}
256 Err(error) => worker::console_error!("could not remove old inbox items: {error}"),
257 }
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo258 match audit::purge(&db, &audit::keep_from(now, max_days), 20).await {
259 Ok(removed) if removed > 0 => {
260 worker::console_log!("removed {removed} audit entries older than {max_days} days")
261 }
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put262 Ok(_) => {}
263 Err(error) => worker::console_error!("could not remove old audit entries: {error}"),
264 }
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo265 // Without billing nobody's plan is known, so nothing younger than the
266 // ceiling is removed.
267 let billing = match env.service("BILLING") {
268 Ok(billing) => billing,
269 Err(error) => {
270 worker::console_error!(
271 "audit entries kept past their plan's days: no billing: {error}"
272 );
273 return;
274 }
275 };
276 match audit::purge_by_plan(&db, &billing, now, min_days, max_days).await {
277 Ok(removed) if removed > 0 => {
278 worker::console_log!("removed {removed} audit entries older than their plan keeps")
279 }
280 Ok(_) => {}
281 Err(error) => worker::console_error!("audit entries kept past their plan's days: {error}"),
282 }
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put283}