pr_01m47d15m3e54sn21z27rpy5n9/services/identity/src/lib.rs

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