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