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/crates/kit/src/lib.rs

520 lines19,190 bytesCodeBlame
1//! Plumbing shared by g1t services that run on Workers.
2//!
3//! Services talk to each other over service bindings with a small JSON
4//! protocol: `POST /rpc/<method>` with the method's arguments as the body,
5//! answered with the method's return value.
6
7use serde::Serialize;
8use serde::de::DeserializeOwned;
9use worker::{Date, Fetcher, Headers, Method, Request, RequestInit, Response, Result};
10
11/// The current time in milliseconds since the epoch.
12pub fn now_ms() -> u64 {
13 Date::now().as_millis()
14}
15
16/// The method name of an RPC request, or `None` if it is not one.
17pub fn rpc_method(request: &Request) -> Option<String> {
18 if request.method() != Method::Post {
19 return None;
20 }
21 request
22 .path()
23 .strip_prefix("/rpc/")
24 .map(|method| method.to_owned())
25}
26
27/// Deserializes a method's arguments.
28pub fn args<A: DeserializeOwned>(body: serde_json::Value) -> Result<A> {
29 serde_json::from_value(body)
30 .map_err(|error| worker::Error::RustError(format!("bad arguments: {error}")))
31}
32
33/// Serializes a method's return value as the response body.
34pub fn reply<R: Serialize>(value: &R) -> Result<Response> {
35 Response::from_json(value)
36}
37
38/// Calls `method` on another service through its binding.
39pub async fn call<A: Serialize, R: DeserializeOwned>(
40 service: &Fetcher,
41 method: &str,
42 arguments: &A,
43) -> Result<R> {
44 let headers = Headers::new();
45 headers.set("content-type", "application/json")?;
46 let mut init = RequestInit::new();
47 init.with_method(Method::Post)
48 .with_headers(headers)
49 .with_body(Some(serde_json::to_string(arguments)?.into()));
50 // The hostname is ignored; a service binding always reaches its service.
51 let request = Request::new_with_init(&format!("https://service/rpc/{method}"), &init)?;
52 let mut response = service.fetch_request(request).await?;
53 if response.status_code() != 200 {
54 return Err(worker::Error::RustError(format!(
55 "{method} failed with status {}: {}",
56 response.status_code(),
57 response.text().await.unwrap_or_default()
58 )));
59 }
60 response.json().await
61}
62
63pub mod d1;
64pub mod wire;
65
66/// Helpers for bindings that workers-rs has no typed wrapper for, such as
67/// Artifacts and Email Sending. Values cross the boundary as JSON.
68pub mod js {
69 use std::fmt;
70
71 use serde::Serialize;
72 use serde::de::DeserializeOwned;
73 use worker::js_sys::{Array, Function, JSON, Promise, Reflect};
74 use worker::wasm_bindgen::{JsCast, JsValue};
75 use worker::wasm_bindgen_futures::JsFuture;
76 use worker::{Env, Error, Result};
77
78 /// An exception thrown by JavaScript, with its `code` if it had one.
79 #[derive(Debug)]
80 pub struct Thrown {
81 pub code: Option<String>,
82 pub message: String,
83 }
84
85 impl Thrown {
86 fn from_value(value: JsValue) -> Self {
87 let property = |name: &str| {
88 Reflect::get(&value, &name.into())
89 .ok()
90 .and_then(|property| property.as_string())
91 };
92 Thrown {
93 code: property("code"),
94 message: property("message").unwrap_or_else(|| format!("{value:?}")),
95 }
96 }
97
98 pub fn is(&self, code: &str) -> bool {
99 self.code.as_deref() == Some(code)
100 }
101 }
102
103 impl fmt::Display for Thrown {
104 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
105 match &self.code {
106 Some(code) => write!(f, "{code}: {}", self.message),
107 None => f.write_str(&self.message),
108 }
109 }
110 }
111
112 impl From<Thrown> for Error {
113 fn from(thrown: Thrown) -> Self {
114 Error::RustError(thrown.to_string())
115 }
116 }
117
118 /// The binding called `name`, as a raw JavaScript value.
119 pub fn binding(env: &Env, name: &str) -> Result<JsValue> {
120 let value = Reflect::get(env.as_ref(), &name.into()).map_err(Thrown::from_value)?;
121 if value.is_undefined() {
122 return Err(Error::RustError(format!(
123 "binding {name} is not configured"
124 )));
125 }
126 Ok(value)
127 }
128
129 /// Reads a property of a JavaScript object.
130 pub fn get(target: &JsValue, name: &str) -> JsValue {
131 Reflect::get(target, &name.into()).unwrap_or(JsValue::UNDEFINED)
132 }
133
134 /// Sets a property on an object.
135 pub fn set(target: &JsValue, name: &str, value: &JsValue) {
136 let _ = Reflect::set(target, &name.into(), value);
137 }
138
139 pub fn to_js<T: Serialize>(value: &T) -> Result<JsValue> {
140 Ok(JSON::parse(&serde_json::to_string(value)?).map_err(Thrown::from_value)?)
141 }
142
143 pub fn from_js<T: DeserializeOwned>(value: &JsValue) -> Result<T> {
144 let text = if value.is_undefined() {
145 None
146 } else {
147 JSON::stringify(value)
148 .map_err(Thrown::from_value)?
149 .as_string()
150 };
151 let text = text.as_deref().unwrap_or("null");
152 serde_json::from_str(text).map_err(|error| {
153 // Say what arrived; a bare serde error is useless in a log.
154 let seen: String = text.chars().take(300).collect();
155 Error::RustError(format!(
156 "unexpected value from JavaScript ({error}): {seen}"
157 ))
158 })
159 }
160
161 /// Calls `target[method](...args)` and awaits the result if it is a
162 /// thenable. `method` may be a name or a symbol.
163 pub async fn call_key(
164 target: &JsValue,
165 method: &JsValue,
166 args: &[JsValue],
167 ) -> std::result::Result<JsValue, Thrown> {
168 let function: Function = Reflect::get(target, method)
169 .map_err(Thrown::from_value)?
170 .dyn_into()
171 .map_err(|_| Thrown {
172 code: None,
173 message: format!("{method:?} is not a function"),
174 })?;
175 let arguments: Array = args.iter().collect();
176 // An RPC stub treats every property access as a remote method, so
177 // `function.apply(...)` would be sent over the wire as a call to
178 // "apply". Reflect.apply invokes the function without touching it.
179 let returned = Reflect::apply(&function, target, &arguments).map_err(Thrown::from_value)?;
180 // Worker RPC returns its own thenable rather than a Promise, so
181 // resolve whatever came back instead of testing its type.
182 JsFuture::from(Promise::resolve(&returned))
183 .await
184 .map_err(Thrown::from_value)
185 }
186
187 /// Calls `target.method(...args)`; see [`call_key`].
188 pub async fn call(
189 target: &JsValue,
190 method: &str,
191 args: &[JsValue],
192 ) -> std::result::Result<JsValue, Thrown> {
193 call_key(target, &method.into(), args).await
194 }
195}
196
197/// Moving a service's rows when a workspace is renamed.
198pub mod rename {
199 use std::collections::HashMap;
200
201 use g1t_contracts::events::{Event, WorkspaceRenamed};
202 use g1t_contracts::identity::UsernamesArgs;
203 use worker::wasm_bindgen::JsValue;
204 use worker::{D1Database, Env, Result};
205
206 /// Handles `workspace.renamed` with `statements`, and says whether
207 /// `event` was one. Each statement uses `?1` for the workspace's current
208 /// slug (asked of identity by id, so renames delivered twice or out of
209 /// order converge) and `?2` for a slug its rows may still be under; the
210 /// statements run in one batch per such slug. A statement that matches
211 /// nothing changes nothing, so running them again is harmless.
212 pub async fn on_event(env: &Env, db: &D1Database, event: &Event, statements: &[&str]) -> Result<bool> {
213 if event.kind != "workspace.renamed" {
214 return Ok(false);
215 }
216 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
217 worker::console_error!("workspace.renamed {} could not be read", event.id);
218 return Ok(true);
219 };
220 let names: HashMap<String, String> = crate::call(
221 &env.service("IDENTITY")?,
222 "usernames",
223 &UsernamesArgs {
224 ids: vec![renamed.workspace_id.clone()],
225 },
226 )
227 .await?;
228 let current = names
229 .get(&renamed.workspace_id)
230 .cloned()
231 .unwrap_or_else(|| renamed.to.clone());
232 for stale in renamed.stale_slugs(&current) {
233 let values: [JsValue; 2] = [current.as_str().into(), stale.as_str().into()];
234 let mut batch = Vec::with_capacity(statements.len());
235 for sql in statements {
236 batch.push(db.prepare(*sql).bind(&values[..parameters(sql)])?);
237 }
238 db.batch(batch).await?;
239 }
240 Ok(true)
241 }
242
243 /// How many values a statement takes: its highest `?N`.
244 pub fn parameters(sql: &str) -> usize {
245 sql.split('?')
246 .skip(1)
247 .filter_map(|rest| {
248 let digits: String = rest.chars().take_while(char::is_ascii_digit).collect();
249 digits.parse().ok()
250 })
251 .max()
252 .unwrap_or(0)
253 }
254
255 #[cfg(test)]
256 mod tests {
257 use super::parameters;
258
259 #[test]
260 fn counts_numbered_parameters() {
261 assert_eq!(parameters("UPDATE t SET a = ?1 WHERE a = ?2"), 2);
262 assert_eq!(parameters("DELETE FROM t WHERE a = ?2"), 2);
263 assert_eq!(parameters("UPDATE t SET a = ?1"), 1);
264 assert_eq!(parameters("DELETE FROM t"), 0);
265 }
266 }
267}
268
269/// Moving a service's rows when a repository's path changes: transferred
270/// to another workspace (`repo.transferred`) or renamed within its own
271/// (`repo.renamed`). Both are handled the same way, so a service that
272/// follows transfers follows renames too.
273pub mod transfer {
274 use g1t_contracts::events::{Event, RepoRenamed, RepoTransferred};
275 use g1t_contracts::repos::{PathByIdArgs, RepoPath};
276 use worker::wasm_bindgen::JsValue;
277 use worker::{D1Database, Env, Result};
278
279 pub use crate::rename::parameters;
280
281 /// A repository's path change, from either event.
282 #[derive(Clone, Debug, PartialEq, Eq)]
283 pub struct Moved {
284 pub repo_id: String,
285 /// The paths (`namespace/name`) the event names, old then new.
286 pub paths: [String; 2],
287 }
288
289 impl Moved {
290 /// The paths whose rows move to `current`: the two the event
291 /// names, minus `current`.
292 pub fn stale_paths(&self, current: &str) -> Vec<String> {
293 let mut paths: Vec<String> = Vec::new();
294 for path in &self.paths {
295 if path != current && !paths.contains(path) {
296 paths.push(path.clone());
297 }
298 }
299 paths
300 }
301
302 /// Where the event says it went, for when repos does not know.
303 pub fn destination(&self) -> &str {
304 &self.paths[1]
305 }
306 }
307
308 impl From<RepoTransferred> for Moved {
309 fn from(t: RepoTransferred) -> Self {
310 Moved {
311 paths: [format!("{}/{}", t.from, t.name), format!("{}/{}", t.to, t.name)],
312 repo_id: t.repo_id,
313 }
314 }
315 }
316
317 impl From<RepoRenamed> for Moved {
318 fn from(r: RepoRenamed) -> Self {
319 Moved {
320 paths: [format!("{}/{}", r.namespace, r.from), format!("{}/{}", r.namespace, r.to)],
321 repo_id: r.repo_id,
322 }
323 }
324 }
325
326 /// The values the statements are bound with for one stale path, in
327 /// order: `?1` the current path (`namespace/name`), `?2` the stale
328 /// path, `?3` the current workspace, `?4` the stale workspace, `?5` the
329 /// repository's id, `?6` the current name, `?7` the stale name.
330 pub fn values(current: &str, stale: &str, repo_id: &str) -> [String; 7] {
331 let namespace = |path: &str| path.split_once('/').map_or(path, |(ns, _)| ns).to_owned();
332 let name = |path: &str| path.split_once('/').map_or("", |(_, name)| name).to_owned();
333 [
334 current.to_owned(),
335 stale.to_owned(),
336 namespace(current),
337 namespace(stale),
338 repo_id.to_owned(),
339 name(current),
340 name(stale),
341 ]
342 }
343
344 /// Handles `repo.transferred` and `repo.renamed` with `statements`, and
345 /// says whether `event` was one. The repository's current path is asked
346 /// of repos by id, so path changes delivered twice or out of order
347 /// converge; the statements run in one batch per path its rows may
348 /// still be under (see [`values`] for the parameters). A statement that
349 /// matches nothing changes nothing, so running them again is harmless.
350 pub async fn on_event(env: &Env, db: &D1Database, event: &Event, statements: &[&str]) -> Result<bool> {
351 let Some(moved) = read(event) else {
352 return Ok(false);
353 };
354 let current = current_path(env, &moved).await?;
355 for stale in moved.stale_paths(&current) {
356 let values = values(&current, &stale, &moved.repo_id);
357 let values: Vec<JsValue> = values.iter().map(|v| JsValue::from(v.as_str())).collect();
358 let mut batch = Vec::with_capacity(statements.len());
359 for sql in statements {
360 batch.push(db.prepare(*sql).bind(&values[..parameters(sql)])?);
361 }
362 db.batch(batch).await?;
363 }
364 Ok(true)
365 }
366
367 /// The path change `event` announces, if it is one.
368 pub fn read(event: &Event) -> Option<Moved> {
369 let read = match event.kind.as_str() {
370 "repo.transferred" => serde_json::from_value::<RepoTransferred>(event.data.clone()).map(Moved::from),
371 "repo.renamed" => serde_json::from_value::<RepoRenamed>(event.data.clone()).map(Moved::from),
372 _ => return None,
373 };
374 if read.is_err() {
375 worker::console_error!("{} {} could not be read", event.kind, event.id);
376 }
377 read.ok()
378 }
379
380 /// The repository's path now, as `namespace/name`: asked of repos, or
381 /// where the event says it went when repos does not know it.
382 pub async fn current_path(env: &Env, moved: &Moved) -> Result<String> {
383 let path: Option<RepoPath> = crate::call(
384 &env.service("REPOS")?,
385 "path_by_id",
386 &PathByIdArgs {
387 id: moved.repo_id.clone(),
388 },
389 )
390 .await?;
391 Ok(path.map_or_else(
392 || moved.destination().to_owned(),
393 |path| format!("{}/{}", path.namespace, path.name),
394 ))
395 }
396
397 #[cfg(test)]
398 mod tests {
399 use super::*;
400
401 #[test]
402 fn binds_paths_workspaces_names_and_the_id() {
403 assert_eq!(
404 values("flagon-io/g1t", "syntaqx/g1t", "rep_1"),
405 ["flagon-io/g1t", "syntaqx/g1t", "flagon-io", "syntaqx", "rep_1", "g1t", "g1t"].map(String::from)
406 );
407 assert_eq!(values("acme/new", "acme/old", "rep_1")[5..], ["new".to_owned(), "old".to_owned()]);
408 }
409
410 #[test]
411 fn a_rename_and_a_transfer_are_both_moves() {
412 let renamed: Moved = RepoRenamed {
413 repo_id: "rep_1".into(),
414 namespace: "acme".into(),
415 from: "old".into(),
416 to: "new".into(),
417 }
418 .into();
419 assert_eq!(renamed.stale_paths("acme/new"), vec!["acme/old"]);
420 assert_eq!(renamed.destination(), "acme/new");
421 let transferred: Moved = RepoTransferred {
422 repo_id: "rep_1".into(),
423 name: "g1t".into(),
424 from: "a".into(),
425 to: "b".into(),
426 }
427 .into();
428 assert_eq!(transferred.stale_paths("c/g1t"), vec!["a/g1t", "b/g1t"]);
429 }
430 }
431}
432
433/// Reading the repository lifecycle events every service reacts to:
434/// `repo.deleted` (stop and hide; it may come back), `repo.restored`
435/// (start again) and `repo.purged` (drop every row kept by its id).
436pub mod lifecycle {
437 use g1t_contracts::events::{Event, RepoDeleted, RepoPurged, RepoRestored};
438 use worker::{D1Database, Result};
439
440 /// One of the three, read.
441 #[derive(Debug)]
442 pub enum Lifecycle {
443 Deleted(RepoDeleted),
444 Restored(RepoRestored),
445 Purged(RepoPurged),
446 }
447
448 impl Lifecycle {
449 pub fn repo_id(&self) -> &str {
450 match self {
451 Lifecycle::Deleted(e) => &e.repo_id,
452 Lifecycle::Restored(e) => &e.repo_id,
453 Lifecycle::Purged(e) => &e.repo_id,
454 }
455 }
456 }
457
458 /// The lifecycle event `event` is, if it is one.
459 pub fn read(event: &Event) -> Option<Lifecycle> {
460 let data = event.data.clone();
461 let read = match event.kind.as_str() {
462 "repo.deleted" => serde_json::from_value(data).map(Lifecycle::Deleted),
463 "repo.restored" => serde_json::from_value(data).map(Lifecycle::Restored),
464 "repo.purged" => serde_json::from_value(data).map(Lifecycle::Purged),
465 _ => return None,
466 };
467 if read.is_err() {
468 worker::console_error!("{} {} could not be read", event.kind, event.id);
469 }
470 read.ok()
471 }
472
473 /// Handles `repo.purged` with `statements`, each taking the
474 /// repository's id as `?1`, in one batch; says whether `event` was
475 /// one. Running them again changes nothing.
476 pub async fn on_purged(db: &D1Database, event: &Event, statements: &[&str]) -> Result<bool> {
477 let Some(Lifecycle::Purged(purged)) = read(event) else {
478 return Ok(false);
479 };
480 if statements.is_empty() {
481 return Ok(true);
482 }
483 let mut batch = Vec::with_capacity(statements.len());
484 for sql in statements {
485 batch.push(db.prepare(*sql).bind(&[purged.repo_id.as_str().into()])?);
486 }
487 db.batch(batch).await?;
488 Ok(true)
489 }
490}
491
492/// Dropping what a service keeps for a workspace alone when the workspace
493/// is deleted.
494pub mod deleted {
495 use g1t_contracts::events::{Event, WorkspaceDeleted};
496 use worker::{D1Database, Result};
497
498 /// Handles `workspace.deleted` with `statements`, each taking the
499 /// workspace's slug as `?1`, in one batch; says whether `event` was
500 /// one. Running them again changes nothing.
501 pub async fn on_event(db: &D1Database, event: &Event, statements: &[&str]) -> Result<bool> {
502 if event.kind != "workspace.deleted" {
503 return Ok(false);
504 }
505 let Ok(deleted) = serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) else {
506 worker::console_error!("workspace.deleted {} could not be read", event.id);
507 return Ok(true);
508 };
509 let slug = deleted.slug.to_lowercase();
510 if slug.is_empty() || statements.is_empty() {
511 return Ok(true);
512 }
513 let mut batch = Vec::with_capacity(statements.len());
514 for sql in statements {
515 batch.push(db.prepare(*sql).bind(&[slug.as_str().into()])?);
516 }
517 db.batch(batch).await?;
518 Ok(true)
519 }
520}