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

644 lines24,680 bytesCodeBlame
1//! The identity service: accounts, sessions, SSH keys and access tokens.
2//!
3//! Reached only through service bindings; see `g1t_contracts::identity` for
4//! the methods and their arguments.
5
6mod admin;
7mod avatars;
8mod crypto;
9mod device;
10mod directory;
11mod email;
12mod oauth;
13mod profiles;
14mod rename;
15mod run_credentials;
16mod tokens;
17mod workspaces;
18
19use g1t_contracts::identity::*;
20use g1t_contracts::time::{SQL_NOW, rfc3339, sql_after};
21use g1t_contracts::{FailureCode, Outcome, User, Viewer, is_valid_namespace, new_id};
22use g1t_kit::{args, now_ms, reply, rpc_method};
23use serde::Deserialize;
24use tokens::TOKEN_PREFIX;
25use worker::wasm_bindgen::JsValue;
26use worker::{Context, D1Database, Env, Request, Response, Result, event};
27
28const SESSION_TTL_SECONDS: u64 = 30 * 24 * 60 * 60;
29const VERIFY_TTL_SECONDS: u64 = 24 * 60 * 60;
30const RESET_TTL_SECONDS: u64 = 60 * 60;
31const MIN_PASSWORD_LENGTH: usize = 10;
32const PASSWORD_TOO_SHORT: &str = "Use a password of at least 10 characters.";
33
34/// A user as selected from the database; `verified` arrives as 0 or 1.
35#[derive(Deserialize)]
36struct Account {
37 id: String,
38 username: String,
39 verified: u8,
40 /// Selected only where the person is being shown to themselves.
41 #[serde(default)]
42 avatar: Option<String>,
43}
44
45impl From<Account> for User {
46 fn from(row: Account) -> Self {
47 User {
48 id: row.id,
49 username: row.username,
50 verified: row.verified != 0,
51 avatar: row.avatar,
52 ..User::default()
53 }
54 }
55}
56
57#[derive(Deserialize)]
58struct UserRow {
59 id: String,
60 username: String,
61 password_hash: String,
62 verified: u8,
63}
64
65/// The owner of an emailed token.
66#[derive(Deserialize)]
67struct TokenOwner {
68 id: String,
69 username: String,
70 email: Option<String>,
71}
72
73#[derive(Deserialize)]
74struct KeyRow {
75 id: String,
76 title: String,
77 fingerprint: String,
78 created_at: String,
79}
80
81impl From<KeyRow> for SshKey {
82 fn from(row: KeyRow) -> Self {
83 SshKey {
84 id: row.id,
85 title: row.title,
86 fingerprint: row.fingerprint,
87 created_at: row.created_at,
88 }
89 }
90}
91
92struct Identity {
93 db: D1Database,
94 env: Env,
95}
96
97impl Identity {
98 /// Runs a query that returns at most one user, for showing to others:
99 /// without their workspaces.
100 async fn find_public_user(&self, sql: &str, param: &str) -> Result<Viewer> {
101 Ok(self
102 .db
103 .prepare(sql)
104 .bind(&[JsValue::from(param)])?
105 .first::<Account>(None)
106 .await?
107 .map(User::from))
108 }
109
110 /// Attaches the workspaces a user belongs to, so that any service can
111 /// authorize them without asking again.
112 async fn with_workspaces(&self, user: Viewer) -> Result<Viewer> {
113 let Some(mut user) = user else {
114 return Ok(None);
115 };
116 user.workspaces = self.memberships(&user.id).await?;
117 Ok(Some(user))
118 }
119
120 /// Runs a query that resolves credentials to at most one user.
121 async fn find_user(&self, sql: &str, param: &str) -> Result<Viewer> {
122 let user = self.find_public_user(sql, param).await?;
123 self.with_workspaces(user).await
124 }
125
126 /// Stores a one-time token of `kind` for the user and returns it.
127 async fn issue_email_token(&self, user_id: &str, kind: &str, ttl: u64) -> Result<String> {
128 let token = crypto::random_hex(32);
129 self.db
130 .prepare(format!(
131 "INSERT INTO email_tokens (id, user_id, kind, expires_at)
132 VALUES (?, ?, ?, {})",
133 sql_after(ttl)
134 ))
135 .bind(&[
136 crypto::sha256_hex(&token).into(),
137 user_id.into(),
138 kind.into(),
139 ])?
140 .run()
141 .await?;
142 Ok(token)
143 }
144
145 /// Consumes a token of `kind`, returning its owner if it was valid.
146 async fn redeem_email_token(&self, token: &str, kind: &str) -> Result<Option<TokenOwner>> {
147 let id = crypto::sha256_hex(token);
148 let owner = self
149 .db
150 .prepare(format!(
151 "SELECT users.id, users.username, users.email FROM email_tokens
152 JOIN users ON users.id = email_tokens.user_id
153 WHERE email_tokens.id = ? AND email_tokens.kind = ?
154 AND email_tokens.expires_at > {SQL_NOW}"
155 ))
156 .bind(&[id.as_str().into(), kind.into()])?
157 .first::<TokenOwner>(None)
158 .await?;
159 if let Some(owner) = &owner {
160 // Every outstanding token of this kind dies with the one used.
161 self.db
162 .prepare("DELETE FROM email_tokens WHERE user_id = ? AND kind = ?")
163 .bind(&[owner.id.as_str().into(), kind.into()])?
164 .run()
165 .await?;
166 }
167 Ok(owner)
168 }
169
170 async fn send_verification(&self, user: &User, email: &str) -> Result<()> {
171 let token = self
172 .issue_email_token(&user.id, "verify", VERIFY_TTL_SECONDS)
173 .await?;
174 email::send_verification(&self.env, email, &user.username, &token).await
175 }
176
177 async fn resend_verification(&self, a: UserArgs) -> Result<Outcome<bool>> {
178 let row = self
179 .db
180 .prepare(
181 "SELECT id, username, email FROM users WHERE id = ? AND email_verified_at IS NULL",
182 )
183 .bind(&[a.user.id.as_str().into()])?
184 .first::<TokenOwner>(None)
185 .await?;
186 let Some(TokenOwner {
187 email: Some(email), ..
188 }) = row
189 else {
190 return Ok(Outcome::fail(
191 FailureCode::Conflict,
192 "This account's email is already confirmed.",
193 ));
194 };
195 self.send_verification(&a.user, &email).await?;
196 Ok(Outcome::Ok(true))
197 }
198
199 async fn verify_email(&self, a: EmailTokenArgs) -> Result<Outcome<User>> {
200 let Some(owner) = self.redeem_email_token(&a.token, "verify").await? else {
201 return Ok(Outcome::fail(
202 FailureCode::Invalid,
203 "This confirmation link is not valid or has expired.",
204 ));
205 };
206 self.db
207 .prepare(format!(
208 "UPDATE users SET email_verified_at = {SQL_NOW} WHERE id = ?"
209 ))
210 .bind(&[owner.id.as_str().into()])?
211 .run()
212 .await?;
213 Ok(Outcome::Ok(User {
214 id: owner.id,
215 username: owner.username,
216 verified: true,
217 ..User::default()
218 }))
219 }
220
221 async fn request_password_reset(&self, a: EmailArgs) -> Result<bool> {
222 let row = self
223 .db
224 .prepare("SELECT id, username, email FROM users WHERE email = ?")
225 .bind(&[a.email.trim().to_lowercase().into()])?
226 .first::<TokenOwner>(None)
227 .await?;
228 if let Some(TokenOwner {
229 id,
230 username,
231 email: Some(email),
232 }) = row
233 {
234 let token = self
235 .issue_email_token(&id, "reset", RESET_TTL_SECONDS)
236 .await?;
237 email::send_password_reset(&self.env, &email, &username, &token).await?;
238 }
239 // The same answer either way, so addresses cannot be probed.
240 Ok(true)
241 }
242
243 async fn reset_password(&self, a: ResetPasswordArgs) -> Result<Outcome<User>> {
244 if a.password.chars().count() < MIN_PASSWORD_LENGTH {
245 return Ok(Outcome::fail(FailureCode::Invalid, PASSWORD_TOO_SHORT));
246 }
247 let Some(owner) = self.redeem_email_token(&a.token, "reset").await? else {
248 return Ok(Outcome::fail(
249 FailureCode::Invalid,
250 "This reset link is not valid or has expired.",
251 ));
252 };
253 // Following an emailed link also proves the address.
254 self.db
255 .prepare(format!(
256 "UPDATE users SET password_hash = ?,
257 email_verified_at = COALESCE(email_verified_at, {SQL_NOW})
258 WHERE id = ?"
259 ))
260 .bind(&[
261 crypto::hash_password(&a.password).into(),
262 owner.id.as_str().into(),
263 ])?
264 .run()
265 .await?;
266 // Anyone signed in with the old password is signed out.
267 self.db
268 .prepare("DELETE FROM sessions WHERE user_id = ?")
269 .bind(&[owner.id.as_str().into()])?
270 .run()
271 .await?;
272 Ok(Outcome::Ok(User {
273 id: owner.id,
274 username: owner.username,
275 verified: true,
276 ..User::default()
277 }))
278 }
279
280 async fn user_for_password(&self, username: &str, password: &str) -> Result<Viewer> {
281 let row = self
282 .db
283 .prepare("SELECT id, username, password_hash, email_verified_at IS NOT NULL AS verified FROM users WHERE username = ?")
284 .bind(&[JsValue::from(username.to_lowercase())])?
285 .first::<UserRow>(None)
286 .await?;
287 let user = row
288 .filter(|row| crypto::verify_password(password, &row.password_hash))
289 .map(|row| User {
290 id: row.id,
291 username: row.username,
292 verified: row.verified != 0,
293 ..User::default()
294 });
295 self.with_workspaces(user).await
296 }
297
298 async fn register(&self, a: RegisterArgs) -> Result<Outcome<SignedIn>> {
299 let username = a.username.trim().to_lowercase();
300 let email = a.email.trim().to_lowercase();
301 let invalid = |message: &str| Ok(Outcome::fail(FailureCode::Invalid, message));
302 if !is_valid_namespace(&username) {
303 return invalid(
304 "Usernames use lowercase letters, digits and single hyphens, up to 39 characters.",
305 );
306 }
307 let well_formed_email = email
308 .split_once('@')
309 .is_some_and(|(local, domain)| !local.is_empty() && domain.contains('.'))
310 && !email.contains(char::is_whitespace);
311 if !well_formed_email {
312 return invalid("Enter a valid email address.");
313 }
314 if a.password.chars().count() < MIN_PASSWORD_LENGTH {
315 return invalid(PASSWORD_TOO_SHORT);
316 }
317 let taken = self
318 .db
319 // Usernames and workspaces share one namespace, so that a name
320 // means the same thing wherever it appears.
321 .prepare(
322 "SELECT username FROM users WHERE username = ? OR email = ?
323 UNION ALL SELECT slug FROM workspaces WHERE slug = ?",
324 )
325 .bind(&[
326 username.as_str().into(),
327 email.as_str().into(),
328 username.as_str().into(),
329 ])?
330 .first::<serde_json::Value>(None)
331 .await?;
332 // A renamed workspace's old slug stays reserved for it a while.
333 if taken.is_some() || self.slug_held(&username).await? {
334 return Ok(Outcome::fail(
335 FailureCode::Conflict,
336 "That username or email is already registered.",
337 ));
338 }
339 let user = User {
340 id: new_id("usr", now_ms()),
341 username,
342 ..User::default()
343 };
344 self.db
345 .prepare("INSERT INTO users (id, username, email, password_hash) VALUES (?, ?, ?, ?)")
346 .bind(&[
347 user.id.as_str().into(),
348 user.username.as_str().into(),
349 email.as_str().into(),
350 crypto::hash_password(&a.password).into(),
351 ])?
352 .run()
353 .await?;
354 // The account exists either way; the email can be sent again later.
355 if let Err(error) = self.send_verification(&user, &email).await {
356 worker::console_error!("verification email failed: {error}");
357 }
358 self.start_session(user).await
359 }
360
361 async fn sign_in(&self, a: SignInArgs) -> Result<Outcome<SignedIn>> {
362 let Some(user) = self.user_for_password(&a.username, &a.password).await? else {
363 return Ok(Outcome::fail(
364 FailureCode::Unauthenticated,
365 "Incorrect username or password.",
366 ));
367 };
368 self.start_session(user).await
369 }
370
371 async fn start_session(&self, user: User) -> Result<Outcome<SignedIn>> {
372 let session_token = crypto::random_hex(32);
373 self.db
374 .prepare(format!(
375 "INSERT INTO sessions (id, user_id, expires_at) VALUES (?, ?, {})",
376 sql_after(SESSION_TTL_SECONDS)
377 ))
378 .bind(&[
379 crypto::sha256_hex(&session_token).into(),
380 user.id.as_str().into(),
381 ])?
382 .run()
383 .await?;
384 Ok(Outcome::Ok(SignedIn {
385 user,
386 session_token,
387 }))
388 }
389
390 async fn sign_out(&self, a: SessionArgs) -> Result<()> {
391 self.db
392 .prepare("DELETE FROM sessions WHERE id = ?")
393 .bind(&[crypto::sha256_hex(&a.session_token).into()])?
394 .run()
395 .await?;
396 Ok(())
397 }
398
399 async fn user_for_session(&self, a: SessionArgs) -> Result<Viewer> {
400 self.find_user(
401 &format!(
402 "SELECT users.id, users.username, users.email_verified_at IS NOT NULL AS verified,
403 users.avatar
404 FROM sessions JOIN users ON users.id = sessions.user_id
405 WHERE sessions.id = ? AND sessions.expires_at > {SQL_NOW}"
406 ),
407 &crypto::sha256_hex(&a.session_token),
408 )
409 .await
410 }
411
412 async fn user_for_git_credentials(&self, a: GitCredentialsArgs) -> Result<Viewer> {
413 // Like GitHub, a token alone identifies its user.
414 if a.secret.starts_with(TOKEN_PREFIX) {
415 self.user_for_access_token(&a.secret).await
416 } else {
417 self.user_for_password(&a.username, &a.secret).await
418 }
419 }
420
421 async fn user_for_ssh_key(&self, a: FingerprintArgs) -> Result<Viewer> {
422 self.find_user(
423 "SELECT users.id, users.username, users.email_verified_at IS NOT NULL AS verified FROM ssh_keys
424 JOIN users ON users.id = ssh_keys.user_id
425 WHERE fingerprint = ?",
426 &a.fingerprint,
427 )
428 .await
429 }
430
431 async fn user_by_username(&self, a: UsernameArgs) -> Result<Viewer> {
432 self.find_public_user(
433 "SELECT id, username, email_verified_at IS NOT NULL AS verified FROM users WHERE username = ?",
434 &a.username.to_lowercase(),
435 )
436 .await
437 }
438
439 async fn usernames(&self, a: UsernamesArgs) -> Result<std::collections::HashMap<String, String>> {
440 #[derive(serde::Deserialize)]
441 struct Named {
442 id: String,
443 name: String,
444 }
445 let ids: Vec<String> = a.ids.into_iter().take(200).collect();
446 let mut names = std::collections::HashMap::new();
447 if ids.is_empty() {
448 return Ok(names);
449 }
450 let marks = vec!["?"; ids.len()].join(", ");
451 let bind: Vec<worker::wasm_bindgen::JsValue> = ids.iter().map(|id| id.as_str().into()).collect();
452 for sql in [
453 format!("SELECT id, username AS name FROM users WHERE id IN ({marks})"),
454 format!("SELECT id, slug AS name FROM workspaces WHERE id IN ({marks})"),
455 ] {
456 for row in self.db.prepare(sql).bind(&bind)?.all().await?.results::<Named>()? {
457 names.insert(row.id, row.name);
458 }
459 }
460 Ok(names)
461 }
462
463 async fn list_ssh_keys(&self, a: UserArgs) -> Result<Vec<SshKey>> {
464 let rows = self
465 .db
466 .prepare("SELECT id, title, fingerprint, created_at FROM ssh_keys WHERE user_id = ? ORDER BY id")
467 .bind(&[a.user.id.into()])?
468 .all()
469 .await?
470 .results::<KeyRow>()?;
471 Ok(rows.into_iter().map(SshKey::from).collect())
472 }
473
474 async fn add_ssh_key(&self, a: AddSshKeyArgs) -> Result<Outcome<SshKey>> {
475 let Some(key) = crypto::parse_ssh_key(&a.public_key) else {
476 return Ok(Outcome::fail(
477 FailureCode::Invalid,
478 "That is not a valid OpenSSH public key.",
479 ));
480 };
481 let taken = self
482 .db
483 .prepare("SELECT id FROM ssh_keys WHERE fingerprint = ?")
484 .bind(&[key.fingerprint.as_str().into()])?
485 .first::<serde_json::Value>(None)
486 .await?;
487 if taken.is_some() {
488 return Ok(Outcome::fail(
489 FailureCode::Conflict,
490 "That key is already registered.",
491 ));
492 }
493 let now = now_ms();
494 let title = [a.title.trim(), key.comment.as_str(), "SSH key"]
495 .into_iter()
496 .find(|candidate| !candidate.is_empty())
497 .unwrap_or_default()
498 .to_owned();
499 let row = KeyRow {
500 id: new_id("key", now),
501 title,
502 fingerprint: key.fingerprint,
503 created_at: rfc3339(now),
504 };
505 self.db
506 .prepare(
507 "INSERT INTO ssh_keys (id, user_id, title, public_key, fingerprint, created_at)
508 VALUES (?, ?, ?, ?, ?, ?)",
509 )
510 .bind(&[
511 row.id.as_str().into(),
512 a.user.id.into(),
513 row.title.as_str().into(),
514 key.public_key.into(),
515 row.fingerprint.as_str().into(),
516 row.created_at.as_str().into(),
517 ])?
518 .run()
519 .await?;
520 Ok(Outcome::Ok(row.into()))
521 }
522
523 /// Deletes a row the user owns from `table`.
524 async fn remove(&self, table: &str, a: RemoveArgs) -> Result<()> {
525 self.db
526 .prepare(format!("DELETE FROM {table} WHERE id = ? AND user_id = ?"))
527 .bind(&[a.id.into(), a.user.id.into()])?
528 .run()
529 .await?;
530 Ok(())
531 }
532}
533
534#[event(fetch)]
535async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
536 let Some(method) = rpc_method(&request) else {
537 return Response::error("Not found", 404);
538 };
539 let body: serde_json::Value = request.json().await?;
540 let identity = Identity {
541 db: env.d1("DB")?,
542 env,
543 };
544
545 match method.as_str() {
546 "register" => {
547 let outcome = identity.register(args(body)?).await?;
548 if let Outcome::Ok(signed_in) = &outcome {
549 identity.announce_user(&signed_in.user.username, Some(&signed_in.user.id)).await;
550 }
551 reply(&outcome)
552 }
553 "sign_in" => reply(&identity.sign_in(args(body)?).await?),
554 "create_workspace" => {
555 let outcome = identity.create_workspace(args(body)?).await?;
556 if let Outcome::Ok(workspace) = &outcome {
557 identity.announce_workspace(&workspace.id, &workspace.slug, None).await;
558 }
559 reply(&outcome)
560 }
561 "get_workspace" => reply(&identity.get_workspace(args(body)?).await?),
562 "list_members" => reply(&identity.list_members(args(body)?).await?),
563 "add_member" => reply(&identity.add_member(args(body)?).await?),
564 "remove_member" => reply(&identity.remove_member(args(body)?).await?),
565 "update_workspace" => {
566 let outcome = identity.update_workspace(args(body)?).await?;
567 if let Outcome::Ok(workspace) = &outcome {
568 identity.announce_workspace(&workspace.id, &workspace.slug, None).await;
569 }
570 reply(&outcome)
571 }
572 "rename_workspace" => reply(&identity.rename_workspace(args(body)?).await?),
573 "check_workspace_rename" => reply(&identity.check_workspace_rename(args(body)?).await?),
574 "resolve_slug" => reply(&identity.resolve_slug(args(body)?).await?),
575 "set_workspace_avatar" => {
576 let outcome = identity.set_workspace_avatar(args(body)?).await?;
577 if let Outcome::Ok(workspace) = &outcome {
578 identity.announce_workspace(&workspace.id, &workspace.slug, None).await;
579 }
580 reply(&outcome)
581 }
582 "set_user_avatar" => {
583 let a: SetUserAvatarArgs = args(body)?;
584 let (username, id) = (a.user.username.clone(), a.user.id.clone());
585 let outcome = identity.set_user_avatar(a).await?;
586 if matches!(outcome, Outcome::Ok(_)) {
587 identity.announce_user(&username, Some(&id)).await;
588 }
589 reply(&outcome)
590 }
591 "list_workspace_tokens" => reply(&identity.list_workspace_tokens(args(body)?).await?),
592 "create_workspace_token" => reply(&identity.create_workspace_token(args(body)?).await?),
593 "remove_workspace_token" => reply(&identity.remove_workspace_token(args(body)?).await?),
594 "oauth_authorize" => reply(&identity.oauth_authorize(args(body)?).await?),
595 "oauth_exchange" => reply(&identity.oauth_exchange(args(body)?).await?),
596 "oauth_refresh" => reply(&identity.oauth_refresh(args(body)?).await?),
597 "list_oauth_grants" => reply(&identity.list_oauth_grants(args(body)?).await?),
598 "revoke_oauth_grant" => reply(&identity.revoke_oauth_grant(args(body)?).await?),
599 "device_start" => reply(&identity.device_start(args(body)?).await?),
600 "device_lookup" => reply(&identity.device_lookup(args(body)?).await?),
601 "device_resolve" => reply(&identity.device_resolve(args(body)?).await?),
602 "device_claim" => reply(&identity.device_claim(args(body)?).await?),
603 "resend_verification" => reply(&identity.resend_verification(args(body)?).await?),
604 "verify_email" => reply(&identity.verify_email(args(body)?).await?),
605 "request_password_reset" => reply(&identity.request_password_reset(args(body)?).await?),
606 "reset_password" => reply(&identity.reset_password(args(body)?).await?),
607 "sign_out" => reply(&identity.sign_out(args(body)?).await?),
608 "user_for_session" => reply(&identity.user_for_session(args(body)?).await?),
609 "user_for_git_credentials" => reply(&identity.user_for_git_credentials(args(body)?).await?),
610 "user_for_access_token" => {
611 let a: TokenArgs = args(body)?;
612 reply(&identity.user_for_access_token(&a.token).await?)
613 }
614 "user_for_ssh_key" => reply(&identity.user_for_ssh_key(args(body)?).await?),
615 "user_by_username" => reply(&identity.user_by_username(args(body)?).await?),
616 "usernames" => reply(&identity.usernames(args(body)?).await?),
617 "profile" => reply(&identity.profile(args(body)?).await?),
618 "update_profile" => {
619 let outcome = identity.update_profile(args(body)?).await?;
620 if let Outcome::Ok(profile) = &outcome {
621 identity.announce_user(&profile.username, None).await;
622 }
623 reply(&outcome)
624 }
625 "directory" => reply(&identity.directory(args(body)?).await?),
626 "profile_workspaces" => reply(&identity.profile_workspaces(args(body)?).await?),
627 "list_ssh_keys" => reply(&identity.list_ssh_keys(args(body)?).await?),
628 "add_ssh_key" => reply(&identity.add_ssh_key(args(body)?).await?),
629 "remove_ssh_key" => reply(&identity.remove("ssh_keys", args(body)?).await?),
630 "list_access_tokens" => reply(&identity.list_access_tokens(args(body)?).await?),
631 "create_access_token" => reply(&identity.create_access_token(args(body)?).await?),
632 "create_agent_token" => reply(&identity.create_agent_token(args(body)?).await?),
633 "agent_scope" => reply(&identity.agent_scope(args(body)?).await?),
634 "create_run_credential" => reply(&identity.create_run_credential(args(body)?).await?),
635 "bind_run_credentials" => reply(&identity.bind_run_credentials(args(body)?).await?),
636 "revoke_run_credentials" => reply(&identity.revoke_run_credentials(args(body)?).await?),
637 "remove_access_token" => reply(&identity.remove("access_tokens", args(body)?).await?),
638 // Staff only: sudo.g1t.sh, over its service binding. See admin.rs.
639 "notify_owners" => reply(&identity.notify_owners(args(body)?).await?),
640 "admin_workspaces" => reply(&identity.admin_workspaces(args(body)?).await?),
641 "admin_workspace" => reply(&identity.admin_workspace(args(body)?).await?),
642 _ => Response::error("Unknown method", 404),
643 }
644}