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

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