flagon-io/g1t

public

Git for AI scale: a forge for thousands of agents working on the same code at once.

g1t/services/events/src/lib.rs

228 lines7,942 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//!
8//! Other services reach it over `POST /rpc/<method>`; see
9//! `g1t_contracts::events` for the methods and their arguments.
10
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API11mod audit;
12
Events service in Rust, with RFC 3339 times and accurate push events13use g1t_contracts::events::{Event, ListArgs, PublishArgs};
14use g1t_contracts::new_id;
15use g1t_contracts::time::rfc3339;
16use g1t_kit::{args, js, now_ms, reply, rpc_method};
17use serde::Deserialize;
18use worker::js_sys::{Array, Object};
19use worker::wasm_bindgen::{JsCast, JsValue};
20use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event};
21
22const DEFAULT_PAGE: u32 = 50;
23const MAX_PAGE: u32 = 200;
24/// Every binding whose name starts with this is a queue that receives all
25/// events: one per subscribing service.
26const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_";
27
28#[derive(Deserialize)]
29struct EventRow {
30 id: String,
31 #[serde(rename = "type")]
32 kind: String,
33 source: String,
34 time: String,
35 repo_id: Option<String>,
36 actor: Option<String>,
37 /// JSON.
38 data: String,
39}
40
41impl From<EventRow> for Event {
42 fn from(row: EventRow) -> Self {
43 Event {
44 id: row.id,
45 kind: row.kind,
46 source: row.source,
47 time: row.time,
48 repo_id: row.repo_id,
49 actor: row.actor,
50 data: serde_json::from_str(&row.data).unwrap_or_default(),
51 }
52 }
53}
54
55fn optional(value: &Option<String>) -> JsValue {
56 value.as_deref().map_or(JsValue::NULL, JsValue::from)
57}
58
59/// Sends `events` to a queue binding, one message each.
60async fn send(queue: &JsValue, events: &[Event]) -> Result<()> {
61 let messages = Array::new();
62 for event in events {
63 let message = Object::new();
64 js::set(&message, "body", &js::to_js(event)?);
65 messages.push(&message);
66 }
67 js::call(queue, "sendBatch", &[messages.into()]).await?;
68 Ok(())
69}
70
71struct Events {
72 db: D1Database,
73 env: Env,
74}
75
76impl Events {
77 /// Assigns each event its id and time and puts it on the bus.
78 async fn publish(&self, a: PublishArgs) -> Result<()> {
79 if a.events.is_empty() {
80 return Ok(());
81 }
82 let now = now_ms();
83 let events: Vec<Event> = a
84 .events
85 .into_iter()
86 .map(|event| Event {
87 id: new_id("evt", now),
88 kind: event.kind,
89 source: event.source,
90 time: rfc3339(now),
91 repo_id: event.repo_id,
92 actor: event.actor,
93 data: event.data,
94 })
95 .collect();
96 send(&js::binding(&self.env, "BUS")?, &events).await
97 }
98
99 /// Newest first. Callers must have checked that the viewer may see the
100 /// repository asked about.
101 async fn list(&self, a: ListArgs) -> Result<Vec<Event>> {
102 let mut conditions = Vec::new();
103 let mut values: Vec<JsValue> = Vec::new();
104 if let Some(repo_id) = &a.repo_id {
105 conditions.push("repo_id = ?".to_owned());
106 values.push(repo_id.as_str().into());
107 }
108 if !a.types.is_empty() {
109 let marks = vec!["?"; a.types.len()].join(", ");
110 conditions.push(format!("type IN ({marks})"));
111 values.extend(a.types.iter().map(|kind| JsValue::from(kind.as_str())));
112 }
113 if let Some(before) = &a.before {
114 conditions.push("id < ?".to_owned());
115 values.push(before.as_str().into());
116 }
117 let filter = if conditions.is_empty() {
118 String::new()
119 } else {
120 format!("WHERE {}", conditions.join(" AND "))
121 };
122 values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into());
123 let rows = self
124 .db
125 .prepare(format!(
126 "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?"
127 ))
128 .bind(&values)?
129 .all()
130 .await?
131 .results::<EventRow>()?;
132 Ok(rows.into_iter().map(Event::from).collect())
133 }
134
135 /// Writes a batch from the bus to the log, then hands it to every
136 /// subscriber.
137 async fn deliver(&self, events: &[Event]) -> Result<()> {
138 let mut statements = Vec::with_capacity(events.len());
139 for event in events {
140 statements.push(
141 self.db
142 .prepare(
143 // Redelivered batches must not duplicate log rows.
144 "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data)
145 VALUES (?, ?, ?, ?, ?, ?, ?)",
146 )
147 .bind(&[
148 event.id.as_str().into(),
149 event.kind.as_str().into(),
150 event.source.as_str().into(),
151 event.time.as_str().into(),
152 optional(&event.repo_id),
153 optional(&event.actor),
154 serde_json::to_string(&event.data)?.into(),
155 ])?,
156 );
157 }
158 self.db.batch(statements).await?;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API159 audit::follow_renames(&self.db, events).await?;
Events service in Rust, with RFC 3339 times and accurate push events160
161 let bindings: &JsValue = self.env.as_ref();
162 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
163 let Some(name) = name
164 .as_string()
165 .filter(|name| name.starts_with(SUBSCRIBER_PREFIX))
166 else {
167 continue;
168 };
169 send(&js::binding(&self.env, &name)?, events).await?;
170 }
171 Ok(())
172 }
173}
174
175#[event(fetch)]
176async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
177 let Some(method) = rpc_method(&request) else {
178 return Response::error("Not found", 404);
179 };
180 let body: serde_json::Value = request.json().await?;
181 let events = Events {
182 db: env.d1("DB")?,
183 env,
184 };
185 match method.as_str() {
186 "publish" => reply(&events.publish(args(body)?).await?),
187 "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 API188 "audit_record" => reply(&audit::record(&events.db, args(body)?).await?),
189 "audit_list" => reply(&audit::list(&events.db, args(body)?).await?),
Events service in Rust, with RFC 3339 times and accurate push events190 _ => Response::error("Unknown method", 404),
191 }
192}
193
194/// Events from the bus. A batch that fails is retried whole, which is safe
195/// because both the log and subscribers ignore an event they have seen.
196#[event(queue)]
197async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
198 let events = Events {
199 db: env.d1("DB")?,
200 env,
201 };
202 let delivered: Vec<Event> = batch
203 .messages()?
204 .into_iter()
205 .map(|message| message.into_body())
206 .collect();
207 events.deliver(&delivered).await?;
208 batch.ack_all();
209 Ok(())
210}
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put211
212/// Once a day: audit entries older than any plan keeps are removed
213/// (`AUDIT_KEEP_DAYS`, a year by default). Shorter windows, such as 30 days
214/// without the Team plan, are applied where the log is read.
215#[event(scheduled)]
216async fn scheduled(_event: worker::ScheduledEvent, env: Env, _ctx: worker::ScheduleContext) {
217 let keep_days = env
218 .var("AUDIT_KEEP_DAYS")
219 .ok()
220 .and_then(|v| v.to_string().parse().ok())
221 .unwrap_or(audit::DEFAULT_KEEP_DAYS);
222 let Ok(db) = env.d1("DB") else { return };
223 match audit::purge(&db, &audit::keep_from(now_ms(), keep_days), 20).await {
224 Ok(removed) if removed > 0 => worker::console_log!("removed {removed} audit entries older than {keep_days} days"),
225 Ok(_) => {}
226 Err(error) => worker::console_error!("could not remove old audit entries: {error}"),
227 }
228}