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 events | 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 | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 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). | |
| Events service in Rust, with RFC 3339 times and accurate push events | 9 | //! |
| Inbox: the events service tells people what needs them as events arrive | 10 | //! It keeps two more things beside the log: the audit log (audit.rs) and |
| Inbox: threads, reasons, subscriptions and watching | 11 | //! each person's inbox (inbox.rs), written as events arrive, with who |
| 12 | //! follows what (subscriptions.rs). | |
| Inbox: the events service tells people what needs them as events arrive | 13 | //! |
| Events service in Rust, with RFC 3339 times and accurate push events | 14 | //! Other services reach it over `POST /rpc/<method>`; see |
| Inbox: the events service tells people what needs them as events arrive | 15 | //! `g1t_contracts::events`, `audit` and `inbox` for the methods and their |
| 16 | //! arguments. | |
| Events service in Rust, with RFC 3339 times and accurate push events | 17 | |
| Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API | 18 | mod audit; |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 19 | mod fanout; |
| Inbox: the events service tells people what needs them as events arrive | 20 | mod inbox; |
| Inbox: threads, reasons, subscriptions and watching | 21 | mod subscriptions; |
| Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API | 22 | |
| Events service in Rust, with RFC 3339 times and accurate push events | 23 | use g1t_contracts::events::{Event, ListArgs, PublishArgs}; |
| 24 | use g1t_contracts::new_id; | |
| 25 | use g1t_contracts::time::rfc3339; | |
| 26 | use g1t_kit::{args, js, now_ms, reply, rpc_method}; | |
| 27 | use serde::Deserialize; | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 28 | use std::collections::{HashMap, HashSet}; |
| Events service in Rust, with RFC 3339 times and accurate push events | 29 | use worker::js_sys::{Array, Object}; |
| 30 | use worker::wasm_bindgen::{JsCast, JsValue}; | |
| 31 | use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event}; | |
| 32 | ||
| 33 | const DEFAULT_PAGE: u32 = 50; | |
| 34 | const MAX_PAGE: u32 = 200; | |
| Merge checks: statuses and check runs on every commit | 35 | /// Statuses and check runs reported on commits: delivered to webhooks, and |
| 36 | /// left out of a timeline unless asked for by type. | |
| 37 | const 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 | ]; | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 46 | /// Every binding whose name starts with this is a queue that receives |
| 47 | /// events: one per subscribing service, sent the types it routes. | |
| Events service in Rust, with RFC 3339 times and accurate push events | 48 | const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_"; |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 49 | /// How long a record of which queues a batch reached is kept for its retries. |
| 50 | const FANOUT_KEEP_MS: u64 = 24 * 60 * 60 * 1000; | |
| Events service in Rust, with RFC 3339 times and accurate push events | 51 | |
| 52 | #[derive(Deserialize)] | |
| 53 | struct 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 | ||
| 65 | impl 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 | ||
| 79 | fn optional(value: &Option<String>) -> JsValue { | |
| 80 | value.as_deref().map_or(JsValue::NULL, JsValue::from) | |
| 81 | } | |
| 82 | ||
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 83 | #[derive(Deserialize)] |
| 84 | struct 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. | |
| 93 | async 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?; | |
| Events service in Rust, with RFC 3339 times and accurate push events | 102 | } |
| 103 | Ok(()) | |
| 104 | } | |
| 105 | ||
| 106 | struct Events { | |
| 107 | db: D1Database, | |
| 108 | env: Env, | |
| 109 | } | |
| 110 | ||
| 111 | impl 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() | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 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 | |
| Events service in Rust, with RFC 3339 times and accurate push events | 134 | }) |
| 135 | .collect(); | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 136 | let events: Vec<&Event> = events.iter().collect(); |
| Events service in Rust, with RFC 3339 times and accurate push events | 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>> { | |
| Merge leftovers: plan activity filtered in the query, pushes counted by account, ghost during the deletion window, profile time zones (identity 0039) | 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(); | |
| Events service in Rust, with RFC 3339 times and accurate push events | 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 | ||
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 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. | |
| Events service in Rust, with RFC 3339 times and accurate push events | 171 | async fn deliver(&self, events: &[Event]) -> Result<()> { |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 172 | let mut statements = Vec::with_capacity(events.len() + 1); |
| Events service in Rust, with RFC 3339 times and accurate push events | 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 | } | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 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 | } | |
| Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API | 207 | audit::follow_renames(&self.db, events).await?; |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 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?; | |
| Events service in Rust, with RFC 3339 times and accurate push events | 211 | |
| 212 | let bindings: &JsValue = self.env.as_ref(); | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 213 | let mut reached: Vec<String> = Vec::new(); |
| 214 | let mut failed = None; | |
| Events service in Rust, with RFC 3339 times and accurate push events | 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 | }; | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 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); | |
| Events service in Rust, with RFC 3339 times and accurate push events | 241 | } |
| Inbox: the events service tells people what needs them as events arrive | 242 | let (work, repos, identity) = ( |
| 243 | self.env.service("WORK")?, | |
| 244 | self.env.service("REPOS")?, | |
| 245 | self.env.service("IDENTITY")?, | |
| 246 | ); | |
| Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002) | 247 | // Optional: without it the inbox still fills, nobody is told live. |
| 248 | let notify = self.env.service("NOTIFY").ok(); | |
| Inbox: the events service tells people what needs them as events arrive | 249 | let sources = inbox::Sources { |
| 250 | work: &work, | |
| 251 | repos: &repos, | |
| 252 | identity: &identity, | |
| Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002) | 253 | notify: notify.as_ref(), |
| Inbox: the events service tells people what needs them as events arrive | 254 | }; |
| 255 | inbox::deliver(&self.db, &sources, events).await; | |
| Events service in Rust, with RFC 3339 times and accurate push events | 256 | Ok(()) |
| 257 | } | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 258 | |
| 259 | /// Writes down which queues each event of a batch has reached, with | |
| 260 | /// those it had reached before, for the batch's retry. | |
| 261 | async fn note_reached( | |
| 262 | &self, | |
| 263 | events: &[Event], | |
| 264 | sent: &HashMap<String, HashSet<String>>, | |
| 265 | reached: &[String], | |
| 266 | ) -> Result<()> { | |
| 267 | let now = rfc3339(now_ms()); | |
| 268 | let mut statements = Vec::with_capacity(events.len()); | |
| 269 | for event in events { | |
| 270 | let mut bindings: Vec<&str> = reached.iter().map(String::as_str).collect(); | |
| 271 | if let Some(before) = sent.get(&event.id) { | |
| 272 | bindings.extend(before.iter().map(String::as_str)); | |
| 273 | } | |
| 274 | bindings.sort_unstable(); | |
| 275 | bindings.dedup(); | |
| 276 | statements.push( | |
| 277 | self.db | |
| 278 | .prepare( | |
| 279 | "INSERT INTO fanout_sent (event_id, bindings, created_at) VALUES (?, ?, ?) | |
| 280 | ON CONFLICT (event_id) DO UPDATE SET bindings = excluded.bindings", | |
| 281 | ) | |
| 282 | .bind(&[event.id.as_str().into(), serde_json::to_string(&bindings)?.into(), now.as_str().into()])?, | |
| 283 | ); | |
| 284 | } | |
| 285 | self.db.batch(statements).await?; | |
| 286 | Ok(()) | |
| 287 | } | |
| Events service in Rust, with RFC 3339 times and accurate push events | 288 | } |
| 289 | ||
| 290 | #[event(fetch)] | |
| 291 | async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> { | |
| 292 | let Some(method) = rpc_method(&request) else { | |
| 293 | return Response::error("Not found", 404); | |
| 294 | }; | |
| 295 | let body: serde_json::Value = request.json().await?; | |
| 296 | let events = Events { | |
| 297 | db: env.d1("DB")?, | |
| 298 | env, | |
| 299 | }; | |
| 300 | match method.as_str() { | |
| 301 | "publish" => reply(&events.publish(args(body)?).await?), | |
| 302 | "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 API | 303 | "audit_record" => reply(&audit::record(&events.db, args(body)?).await?), |
| 304 | "audit_list" => reply(&audit::list(&events.db, args(body)?).await?), | |
| Inbox: the events service tells people what needs them as events arrive | 305 | "inbox_list" => { |
| 306 | let repos = events.env.service("REPOS")?; | |
| 307 | reply(&inbox::list(&events.db, &repos, args(body)?).await?) | |
| 308 | } | |
| 309 | "inbox_counts" => reply(&inbox::counts(&events.db, args(body)?).await?), | |
| Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002) | 310 | "inbox_mark" => { |
| 311 | let marked: g1t_contracts::inbox::MarkInboxArgs = args(body)?; | |
| 312 | let username = marked.username.clone(); | |
| 313 | let changed = inbox::mark(&events.db, marked).await?; | |
| 314 | // Notify: the new count in every open tab, whoever marked (the site, the API, MCP). | |
| 315 | if changed > 0 | |
| 316 | && let Ok(notify) = events.env.service("NOTIFY") | |
| 317 | { | |
| 318 | inbox::tell_inbox_count(&events.db, ¬ify, &username).await; | |
| 319 | } | |
| 320 | reply(&changed) | |
| 321 | } | |
| Inbox: threads, reasons, subscriptions and watching | 322 | "inbox_thread" => { |
| 323 | let (repos, work) = (events.env.service("REPOS")?, events.env.service("WORK")?); | |
| 324 | reply(&inbox::thread(&events.db, &repos, &work, args(body)?).await?) | |
| 325 | } | |
| 326 | "inbox_subscription" => { | |
| 327 | let work = events.env.service("WORK")?; | |
| 328 | reply(&subscriptions::subscription(&events.db, &work, args(body)?).await?) | |
| 329 | } | |
| 330 | "inbox_subscribe" => { | |
| 331 | let work = events.env.service("WORK")?; | |
| 332 | reply(&subscriptions::subscribe(&events.db, &work, args(body)?).await?) | |
| 333 | } | |
| 334 | "inbox_watching" => reply(&subscriptions::watching(&events.db, args(body)?).await?), | |
| 335 | "inbox_watch" => reply(&subscriptions::watch(&events.db, args(body)?).await?), | |
| Merge branch 'main' into worktree-agent-a69aeabc4b0deeb97 | 336 | "inbox_watchers" => reply(&subscriptions::watchers_count(&events.db, args(body)?).await?), |
| Inbox: threads, reasons, subscriptions and watching | 337 | "inbox_watched" => reply(&subscriptions::watched(&events.db, args(body)?).await?), |
| 338 | "inbox_settings" => reply(&subscriptions::settings(&events.db, args(body)?).await?), | |
| 339 | "inbox_update_settings" => reply(&subscriptions::update_settings(&events.db, args(body)?).await?), | |
| Events service in Rust, with RFC 3339 times and accurate push events | 340 | _ => Response::error("Unknown method", 404), |
| 341 | } | |
| 342 | } | |
| 343 | ||
| 344 | /// Events from the bus. A batch that fails is retried whole, which is safe | |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 345 | /// because the log ignores an event it has seen, and queues that already |
| 346 | /// have it are skipped (`deliver`). After its retries it goes to the | |
| 347 | /// dead-letter queue (wrangler.jsonc). | |
| Events service in Rust, with RFC 3339 times and accurate push events | 348 | #[event(queue)] |
| 349 | async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> { | |
| 350 | let events = Events { | |
| 351 | db: env.d1("DB")?, | |
| 352 | env, | |
| 353 | }; | |
| 354 | let delivered: Vec<Event> = batch | |
| 355 | .messages()? | |
| 356 | .into_iter() | |
| 357 | .map(|message| message.into_body()) | |
| 358 | .collect(); | |
| 359 | events.deliver(&delivered).await?; | |
| 360 | batch.ack_all(); | |
| 361 | Ok(()) | |
| 362 | } | |
| Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put | 363 | |
| Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo | 364 | /// Once a day, old audit entries are removed: first everything older than |
| 365 | /// any workspace keeps (`AUDIT_MAX_DAYS`, 400 days), then each workspace's | |
| 366 | /// entries older than its own plan keeps, as billing says (7 days free, 90 | |
| 367 | /// on the plan, or what staff set). Workspaces with nothing older than the | |
| 368 | /// 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 put | 369 | #[event(scheduled)] |
| 370 | async 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 sudo | 371 | let days = |name: &str, default: u32| { |
| 372 | env.var(name) | |
| 373 | .ok() | |
| 374 | .and_then(|v| v.to_string().trim().parse::<u32>().ok()) | |
| 375 | .filter(|days| *days > 0) | |
| 376 | .unwrap_or(default) | |
| 377 | }; | |
| 378 | let max_days = days("AUDIT_MAX_DAYS", audit::DEFAULT_MAX_DAYS); | |
| 379 | 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 put | 380 | 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 sudo | 381 | let now = now_ms(); |
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 382 | let cutoff = rfc3339(now.saturating_sub(FANOUT_KEEP_MS)); |
| 383 | let purged = match db.prepare("DELETE FROM fanout_sent WHERE created_at < ?").bind(&[cutoff.into()]) { | |
| 384 | Ok(statement) => statement.run().await.map(|_| ()), | |
| 385 | Err(error) => Err(error), | |
| 386 | }; | |
| 387 | if let Err(error) = purged { | |
| 388 | worker::console_error!("could not remove old fan-out records: {error}"); | |
| 389 | } | |
| Inbox: the events service tells people what needs them as events arrive | 390 | match inbox::purge(&db, now).await { |
| 391 | Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"), | |
| 392 | Ok(_) => {} | |
| 393 | Err(error) => worker::console_error!("could not remove old inbox items: {error}"), | |
| 394 | } | |
| Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo | 395 | match audit::purge(&db, &audit::keep_from(now, max_days), 20).await { |
| 396 | Ok(removed) if removed > 0 => { | |
| 397 | worker::console_log!("removed {removed} audit entries older than {max_days} days") | |
| 398 | } | |
| Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put | 399 | Ok(_) => {} |
| 400 | Err(error) => worker::console_error!("could not remove old audit entries: {error}"), | |
| 401 | } | |
| Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo | 402 | // Without billing nobody's plan is known, so nothing younger than the |
| 403 | // ceiling is removed. | |
| 404 | let billing = match env.service("BILLING") { | |
| 405 | Ok(billing) => billing, | |
| 406 | Err(error) => { | |
| 407 | worker::console_error!( | |
| 408 | "audit entries kept past their plan's days: no billing: {error}" | |
| 409 | ); | |
| 410 | return; | |
| 411 | } | |
| 412 | }; | |
| 413 | match audit::purge_by_plan(&db, &billing, now, min_days, max_days).await { | |
| 414 | Ok(removed) if removed > 0 => { | |
| 415 | worker::console_log!("removed {removed} audit entries older than their plan keeps") | |
| 416 | } | |
| 417 | Ok(_) => {} | |
| 418 | Err(error) => worker::console_error!("audit entries kept past their plan's days: {error}"), | |
| 419 | } | |
| Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put | 420 | } |
| Merge leftovers: plan activity filtered in the query, pushes counted by account, ghost during the deletion window, profile time zones (identity 0039) | 421 | |
| 422 | /// A value bound to a `?` in [`list_filter`]'s clause. | |
| 423 | #[derive(Debug, PartialEq)] | |
| 424 | enum Bind { | |
| 425 | Text(String), | |
| 426 | Number(f64), | |
| 427 | } | |
| 428 | ||
| 429 | /// The `WHERE` clause `list` reads with, and what it binds, in order. | |
| 430 | fn list_filter(a: &ListArgs) -> (String, Vec<Bind>) { | |
| 431 | let mut conditions = Vec::new(); | |
| 432 | let mut values = Vec::new(); | |
| 433 | if let Some(repo_id) = &a.repo_id { | |
| 434 | conditions.push("repo_id = ?".to_owned()); | |
| 435 | values.push(Bind::Text(repo_id.clone())); | |
| 436 | } | |
| 437 | if !a.types.is_empty() { | |
| 438 | let marks = vec!["?"; a.types.len()].join(", "); | |
| 439 | conditions.push(format!("type IN ({marks})")); | |
| 440 | values.extend(a.types.iter().map(|kind| Bind::Text(kind.clone()))); | |
| 441 | } else { | |
| 442 | // What CI and integrations report on commits goes to webhooks, | |
| 443 | // and is read from each commit's checks; a timeline asked for | |
| 444 | // everything would be little else on a busy repository. | |
| 445 | let marks = vec!["?"; REPORTING.len()].join(", "); | |
| 446 | conditions.push(format!("type NOT IN ({marks})")); | |
| 447 | values.extend(REPORTING.iter().map(|kind| Bind::Text((*kind).to_owned()))); | |
| 448 | } | |
| 449 | if let Some(actor) = &a.actor { | |
| 450 | conditions.push("actor = ?".to_owned()); | |
| 451 | values.push(Bind::Text(actor.clone())); | |
| 452 | } | |
| 453 | if !a.numbers.is_empty() { | |
| 454 | // An issue or pull request's own events name it as `number`; a | |
| 455 | // comment, review or link on it names it as `issue`. | |
| 456 | let marks = vec!["?"; a.numbers.len()].join(", "); | |
| 457 | conditions.push(format!( | |
| 458 | "(json_extract(data, '$.number') IN ({marks}) OR json_extract(data, '$.issue') IN ({marks}))" | |
| 459 | )); | |
| 460 | for _ in 0..2 { | |
| 461 | values.extend(a.numbers.iter().map(|number| Bind::Number(f64::from(*number)))); | |
| 462 | } | |
| 463 | } | |
| 464 | if let Some(since) = &a.since { | |
| 465 | conditions.push("time >= ?".to_owned()); | |
| 466 | values.push(Bind::Text(since.clone())); | |
| 467 | } | |
| 468 | if let Some(before) = &a.before { | |
| 469 | conditions.push("id < ?".to_owned()); | |
| 470 | values.push(Bind::Text(before.clone())); | |
| 471 | } | |
| 472 | let filter = if conditions.is_empty() { | |
| 473 | String::new() | |
| 474 | } else { | |
| 475 | format!("WHERE {}", conditions.join(" AND ")) | |
| 476 | }; | |
| 477 | (filter, values) | |
| 478 | } | |
| 479 | ||
| 480 | #[cfg(test)] | |
| 481 | mod tests { | |
| 482 | use super::*; | |
| 483 | ||
| 484 | #[test] | |
| 485 | fn a_plain_list_leaves_out_what_is_reported_on_commits() { | |
| 486 | let (filter, binds) = list_filter(&ListArgs { repo_id: Some("rep_1".into()), ..ListArgs::default() }); | |
| 487 | assert!(filter.starts_with("WHERE repo_id = ? AND type NOT IN (")); | |
| 488 | assert_eq!(binds[0], Bind::Text("rep_1".into())); | |
| 489 | assert_eq!(binds.len(), 1 + REPORTING.len()); | |
| 490 | } | |
| 491 | ||
| 492 | #[test] | |
| 493 | fn numbers_match_an_item_or_what_is_said_on_it_since_a_time() { | |
| 494 | let (filter, binds) = list_filter(&ListArgs { | |
| 495 | repo_id: Some("rep_1".into()), | |
| 496 | types: vec!["issue.opened".into()], | |
| 497 | numbers: vec![4, 9], | |
| 498 | since: Some("2026-10-01T00:00:00Z".into()), | |
| 499 | actor: Some("usr_g1t_agent".into()), | |
| 500 | ..ListArgs::default() | |
| 501 | }); | |
| 502 | assert_eq!( | |
| 503 | filter, | |
| 504 | "WHERE repo_id = ? AND type IN (?) AND actor = ? AND \ | |
| 505 | (json_extract(data, '$.number') IN (?, ?) OR json_extract(data, '$.issue') IN (?, ?)) AND time >= ?" | |
| 506 | ); | |
| 507 | assert_eq!( | |
| 508 | binds, | |
| 509 | vec![ | |
| 510 | Bind::Text("rep_1".into()), | |
| 511 | Bind::Text("issue.opened".into()), | |
| 512 | Bind::Text("usr_g1t_agent".into()), | |
| 513 | Bind::Number(4.0), | |
| 514 | Bind::Number(9.0), | |
| 515 | Bind::Number(4.0), | |
| 516 | Bind::Number(9.0), | |
| 517 | Bind::Text("2026-10-01T00:00:00Z".into()), | |
| 518 | ] | |
| 519 | ); | |
| 520 | } | |
| 521 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.