flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/events/src/lib.rs

205 lines6,857 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
11use g1t_contracts::events::{Event, ListArgs, PublishArgs};
12use g1t_contracts::new_id;
13use g1t_contracts::time::rfc3339;
14use g1t_kit::{args, js, now_ms, reply, rpc_method};
15use serde::Deserialize;
16use worker::js_sys::{Array, Object};
17use worker::wasm_bindgen::{JsCast, JsValue};
18use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event};
19
20const DEFAULT_PAGE: u32 = 50;
21const MAX_PAGE: u32 = 200;
22/// Every binding whose name starts with this is a queue that receives all
23/// events: one per subscribing service.
24const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_";
25
26#[derive(Deserialize)]
27struct EventRow {
28 id: String,
29 #[serde(rename = "type")]
30 kind: String,
31 source: String,
32 time: String,
33 repo_id: Option<String>,
34 actor: Option<String>,
35 /// JSON.
36 data: String,
37}
38
39impl From<EventRow> for Event {
40 fn from(row: EventRow) -> Self {
41 Event {
42 id: row.id,
43 kind: row.kind,
44 source: row.source,
45 time: row.time,
46 repo_id: row.repo_id,
47 actor: row.actor,
48 data: serde_json::from_str(&row.data).unwrap_or_default(),
49 }
50 }
51}
52
53fn optional(value: &Option<String>) -> JsValue {
54 value.as_deref().map_or(JsValue::NULL, JsValue::from)
55}
56
57/// Sends `events` to a queue binding, one message each.
58async fn send(queue: &JsValue, events: &[Event]) -> Result<()> {
59 let messages = Array::new();
60 for event in events {
61 let message = Object::new();
62 js::set(&message, "body", &js::to_js(event)?);
63 messages.push(&message);
64 }
65 js::call(queue, "sendBatch", &[messages.into()]).await?;
66 Ok(())
67}
68
69struct Events {
70 db: D1Database,
71 env: Env,
72}
73
74impl Events {
75 /// Assigns each event its id and time and puts it on the bus.
76 async fn publish(&self, a: PublishArgs) -> Result<()> {
77 if a.events.is_empty() {
78 return Ok(());
79 }
80 let now = now_ms();
81 let events: Vec<Event> = a
82 .events
83 .into_iter()
84 .map(|event| Event {
85 id: new_id("evt", now),
86 kind: event.kind,
87 source: event.source,
88 time: rfc3339(now),
89 repo_id: event.repo_id,
90 actor: event.actor,
91 data: event.data,
92 })
93 .collect();
94 send(&js::binding(&self.env, "BUS")?, &events).await
95 }
96
97 /// Newest first. Callers must have checked that the viewer may see the
98 /// repository asked about.
99 async fn list(&self, a: ListArgs) -> Result<Vec<Event>> {
100 let mut conditions = Vec::new();
101 let mut values: Vec<JsValue> = Vec::new();
102 if let Some(repo_id) = &a.repo_id {
103 conditions.push("repo_id = ?".to_owned());
104 values.push(repo_id.as_str().into());
105 }
106 if !a.types.is_empty() {
107 let marks = vec!["?"; a.types.len()].join(", ");
108 conditions.push(format!("type IN ({marks})"));
109 values.extend(a.types.iter().map(|kind| JsValue::from(kind.as_str())));
110 }
111 if let Some(before) = &a.before {
112 conditions.push("id < ?".to_owned());
113 values.push(before.as_str().into());
114 }
115 let filter = if conditions.is_empty() {
116 String::new()
117 } else {
118 format!("WHERE {}", conditions.join(" AND "))
119 };
120 values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into());
121 let rows = self
122 .db
123 .prepare(format!(
124 "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?"
125 ))
126 .bind(&values)?
127 .all()
128 .await?
129 .results::<EventRow>()?;
130 Ok(rows.into_iter().map(Event::from).collect())
131 }
132
133 /// Writes a batch from the bus to the log, then hands it to every
134 /// subscriber.
135 async fn deliver(&self, events: &[Event]) -> Result<()> {
136 let mut statements = Vec::with_capacity(events.len());
137 for event in events {
138 statements.push(
139 self.db
140 .prepare(
141 // Redelivered batches must not duplicate log rows.
142 "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data)
143 VALUES (?, ?, ?, ?, ?, ?, ?)",
144 )
145 .bind(&[
146 event.id.as_str().into(),
147 event.kind.as_str().into(),
148 event.source.as_str().into(),
149 event.time.as_str().into(),
150 optional(&event.repo_id),
151 optional(&event.actor),
152 serde_json::to_string(&event.data)?.into(),
153 ])?,
154 );
155 }
156 self.db.batch(statements).await?;
157
158 let bindings: &JsValue = self.env.as_ref();
159 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
160 let Some(name) = name
161 .as_string()
162 .filter(|name| name.starts_with(SUBSCRIBER_PREFIX))
163 else {
164 continue;
165 };
166 send(&js::binding(&self.env, &name)?, events).await?;
167 }
168 Ok(())
169 }
170}
171
172#[event(fetch)]
173async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
174 let Some(method) = rpc_method(&request) else {
175 return Response::error("Not found", 404);
176 };
177 let body: serde_json::Value = request.json().await?;
178 let events = Events {
179 db: env.d1("DB")?,
180 env,
181 };
182 match method.as_str() {
183 "publish" => reply(&events.publish(args(body)?).await?),
184 "list" => reply(&events.list(args(body)?).await?),
185 _ => Response::error("Unknown method", 404),
186 }
187}
188
189/// Events from the bus. A batch that fails is retried whole, which is safe
190/// because both the log and subscribers ignore an event they have seen.
191#[event(queue)]
192async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
193 let events = Events {
194 db: env.d1("DB")?,
195 env,
196 };
197 let delivered: Vec<Event> = batch
198 .messages()?
199 .into_iter()
200 .map(|message| message.into_body())
201 .collect();
202 events.deliver(&delivered).await?;
203 batch.ack_all();
204 Ok(())
205}