pr_01m47d15m3e54sn21z27rpy5n9/services/identity/src/lib.rs

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