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

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