Skip to content
574 linesCodeBlameRaw

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.

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

This file's history is long; its oldest lines are credited to the oldest commit read.