pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/authored.rs
| 1 | //! A person's work, for their profile: the issues and pull requests they |
| 2 | //! opened, newest first, a page at a time. |
| 3 | //! |
| 4 | //! Only work on repositories the viewer may read is ever shown, counted or |
| 5 | //! named. Which those are is the repos service's decision (`readable`, the |
| 6 | //! same check as opening the repository), made once per request over every |
| 7 | //! repository the person has worked in; every query here is then confined |
| 8 | //! to that set. A private repository's titles, numbers and even its |
| 9 | //! existence never reach anyone who could not open it. |
| 10 | |
| 11 | use std::collections::HashMap; |
| 12 | |
| 13 | use g1t_contracts::identity::UsernameArgs; |
| 14 | use g1t_contracts::repos::{ReadableArgs, Repo, RepoPath, MAX_READABLE}; |
| 15 | use g1t_contracts::work::*; |
| 16 | use g1t_contracts::{FailureCode, Outcome, Viewer}; |
| 17 | use serde::Deserialize; |
| 18 | use worker::Result; |
| 19 | use worker::wasm_bindgen::JsValue; |
| 20 | |
| 21 | use crate::Work; |
| 22 | |
| 23 | #[derive(Deserialize)] |
| 24 | struct ItemRow { |
| 25 | kind: AuthoredKind, |
| 26 | id: String, |
| 27 | repo_id: String, |
| 28 | number: u32, |
| 29 | title: String, |
| 30 | state: State, |
| 31 | reason: Option<IssueReason>, |
| 32 | status: Option<PullStatus>, |
| 33 | created_at: String, |
| 34 | updated_at: String, |
| 35 | merged_at: Option<String>, |
| 36 | } |
| 37 | |
| 38 | #[derive(Deserialize)] |
| 39 | struct RepoCount { |
| 40 | repo_id: String, |
| 41 | n: u32, |
| 42 | } |
| 43 | |
| 44 | /// Where a page ends: the sort key and id of its last item. |
| 45 | fn cursor(sort: AuthoredSort, row: &ItemRow) -> String { |
| 46 | let key = match sort { |
| 47 | AuthoredSort::Updated => &row.updated_at, |
| 48 | AuthoredSort::Created | AuthoredSort::Oldest => &row.created_at, |
| 49 | }; |
| 50 | format!("{key}|{}", row.id) |
| 51 | } |
| 52 | |
| 53 | /// A cursor read back, or none if it is not one. |
| 54 | fn parse_cursor(value: &str) -> Option<(&str, &str)> { |
| 55 | let (key, id) = value.split_once('|')?; |
| 56 | (!key.is_empty() && !id.is_empty() && !id.contains('|')).then_some((key, id)) |
| 57 | } |
| 58 | |
| 59 | /// One bound value; made a `JsValue` only when bound, so the SQL can be |
| 60 | /// built and tested outside a Worker. |
| 61 | enum Param { |
| 62 | Text(String), |
| 63 | Number(u32), |
| 64 | } |
| 65 | |
| 66 | impl From<&str> for Param { |
| 67 | fn from(value: &str) -> Self { |
| 68 | Param::Text(value.to_owned()) |
| 69 | } |
| 70 | } |
| 71 | |
| 72 | impl From<u32> for Param { |
| 73 | fn from(value: u32) -> Self { |
| 74 | Param::Number(value) |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | /// Parameters bound by number, so each can be used more than once. |
| 79 | struct Params(Vec<Param>); |
| 80 | |
| 81 | impl Params { |
| 82 | fn push(&mut self, value: impl Into<Param>) -> String { |
| 83 | self.0.push(value.into()); |
| 84 | format!("?{}", self.0.len()) |
| 85 | } |
| 86 | |
| 87 | fn values(&self) -> Vec<JsValue> { |
| 88 | self.0 |
| 89 | .iter() |
| 90 | .map(|param| match param { |
| 91 | Param::Text(text) => JsValue::from(text.as_str()), |
| 92 | Param::Number(n) => JsValue::from(*n), |
| 93 | }) |
| 94 | .collect() |
| 95 | } |
| 96 | } |
| 97 | |
| 98 | /// The work SQL for one page: the person's issues and pull requests in |
| 99 | /// `visible` as one list, filtered, ordered and cut to `limit + 1` rows. |
| 100 | fn page_sql(a: &ByAuthorArgs, params: &mut Params, author: &str, visible: &str, limit: u32) -> String { |
| 101 | let author = params.push(author); |
| 102 | let visible = params.push(visible); |
| 103 | let mut conditions = Vec::new(); |
| 104 | if let Some(kind) = a.kind { |
| 105 | conditions.push(format!( |
| 106 | "kind = {}", |
| 107 | params.push(match kind { |
| 108 | AuthoredKind::Issue => "issue", |
| 109 | AuthoredKind::Pull => "pull", |
| 110 | }) |
| 111 | )); |
| 112 | } |
| 113 | match a.state { |
| 114 | Some(AuthoredState::Open) => conditions.push("state = 'open'".to_owned()), |
| 115 | Some(AuthoredState::Closed) => conditions.push("state = 'closed'".to_owned()), |
| 116 | Some(AuthoredState::Merged) => conditions.push("status = 'merged'".to_owned()), |
| 117 | None => {} |
| 118 | } |
| 119 | let (key, direction, beyond) = match a.sort { |
| 120 | AuthoredSort::Created => ("created_at", "DESC", "<"), |
| 121 | AuthoredSort::Updated => ("updated_at", "DESC", "<"), |
| 122 | AuthoredSort::Oldest => ("created_at", "ASC", ">"), |
| 123 | }; |
| 124 | if let Some((after_key, after_id)) = a.before.as_deref().and_then(parse_cursor) { |
| 125 | let after_key = params.push(after_key); |
| 126 | let after_id = params.push(after_id); |
| 127 | conditions.push(format!( |
| 128 | "({key} {beyond} {after_key} OR ({key} = {after_key} AND id {beyond} {after_id}))" |
| 129 | )); |
| 130 | } |
| 131 | let filter = if conditions.is_empty() { |
| 132 | String::new() |
| 133 | } else { |
| 134 | format!("WHERE {}", conditions.join(" AND ")) |
| 135 | }; |
| 136 | let limit = params.push(limit + 1); |
| 137 | format!( |
| 138 | "SELECT * FROM ( |
| 139 | SELECT 'issue' AS kind, id, repo_id, number, title, state, reason, |
| 140 | NULL AS status, created_at, updated_at, NULL AS merged_at |
| 141 | FROM issues |
| 142 | WHERE author_id = {author} AND repo_id IN (SELECT value FROM json_each({visible})) |
| 143 | UNION ALL |
| 144 | SELECT 'pull' AS kind, id, repo_id, number, title, |
| 145 | CASE WHEN status IN ('draft', 'open') THEN 'open' ELSE 'closed' END AS state, |
| 146 | NULL AS reason, status, created_at, updated_at, merged_at |
| 147 | FROM pulls |
| 148 | WHERE author_id = {author} AND repo_id IN (SELECT value FROM json_each({visible})) |
| 149 | ) AS work |
| 150 | {filter} |
| 151 | ORDER BY {key} {direction}, id {direction} |
| 152 | LIMIT {limit}" |
| 153 | ) |
| 154 | } |
| 155 | |
| 156 | const COUNTS_SQL: &str = "SELECT |
| 157 | (SELECT count(*) FROM pulls WHERE author_id = ?1 |
| 158 | AND repo_id IN (SELECT value FROM json_each(?2)) AND status = 'merged') AS pulls_merged, |
| 159 | (SELECT count(*) FROM pulls WHERE author_id = ?1 |
| 160 | AND repo_id IN (SELECT value FROM json_each(?2)) AND status IN ('draft', 'open')) AS pulls_open, |
| 161 | (SELECT count(*) FROM pulls WHERE author_id = ?1 |
| 162 | AND repo_id IN (SELECT value FROM json_each(?2))) AS pulls, |
| 163 | (SELECT count(*) FROM issues WHERE author_id = ?1 |
| 164 | AND repo_id IN (SELECT value FROM json_each(?2))) AS issues, |
| 165 | (SELECT count(*) FROM issues WHERE author_id = ?1 |
| 166 | AND repo_id IN (SELECT value FROM json_each(?2)) AND state = 'open') AS issues_open"; |
| 167 | |
| 168 | #[derive(Deserialize)] |
| 169 | #[serde(rename_all = "snake_case")] |
| 170 | struct CountsRow { |
| 171 | pulls_merged: u32, |
| 172 | pulls_open: u32, |
| 173 | pulls: u32, |
| 174 | issues: u32, |
| 175 | issues_open: u32, |
| 176 | } |
| 177 | |
| 178 | /// Whether `repo` is the one named `namespace/name`. |
| 179 | fn is_named(repo: &Repo, name: &str) -> bool { |
| 180 | name.split_once('/').is_some_and(|(namespace, name)| { |
| 181 | repo.namespace.eq_ignore_ascii_case(namespace) && repo.name.eq_ignore_ascii_case(name) |
| 182 | }) |
| 183 | } |
| 184 | |
| 185 | impl Work { |
| 186 | pub(crate) async fn by_author(&self, a: ByAuthorArgs) -> Result<Outcome<Authored>> { |
| 187 | let person: Viewer = g1t_kit::call( |
| 188 | &self.identity, |
| 189 | "user_by_username", |
| 190 | &UsernameArgs { |
| 191 | username: a.username.trim().to_lowercase(), |
| 192 | }, |
| 193 | ) |
| 194 | .await?; |
| 195 | let Some(person) = person else { |
| 196 | return Ok(Outcome::fail(FailureCode::NotFound, "There is no such account.")); |
| 197 | }; |
| 198 | |
| 199 | // Every repository they have opened work in, busiest first... |
| 200 | let touched = self |
| 201 | .db |
| 202 | .prepare( |
| 203 | "SELECT repo_id, count(*) AS n FROM ( |
| 204 | SELECT repo_id FROM issues WHERE author_id = ?1 |
| 205 | UNION ALL SELECT repo_id FROM pulls WHERE author_id = ?1 |
| 206 | ) GROUP BY repo_id ORDER BY n DESC LIMIT ?2", |
| 207 | ) |
| 208 | .bind(&[person.id.as_str().into(), (MAX_READABLE as u32).into()])? |
| 209 | .all() |
| 210 | .await? |
| 211 | .results::<RepoCount>()?; |
| 212 | if touched.is_empty() { |
| 213 | return Ok(Outcome::Ok(Authored::default())); |
| 214 | } |
| 215 | // ...and of those, the ones the viewer may read. |
| 216 | let readable: Vec<Repo> = g1t_kit::call( |
| 217 | &self.repos, |
| 218 | "readable", |
| 219 | &ReadableArgs { |
| 220 | ids: touched.iter().map(|row| row.repo_id.clone()).collect(), |
| 221 | viewer: a.viewer.clone(), |
| 222 | }, |
| 223 | ) |
| 224 | .await?; |
| 225 | let by_id: HashMap<&str, &Repo> = readable.iter().map(|repo| (repo.id.as_str(), repo)).collect(); |
| 226 | let repos: Vec<AuthoredRepo> = touched |
| 227 | .iter() |
| 228 | .filter_map(|row| { |
| 229 | by_id.get(row.repo_id.as_str()).map(|repo| AuthoredRepo { |
| 230 | repo: RepoPath { |
| 231 | namespace: repo.namespace.clone(), |
| 232 | name: repo.name.clone(), |
| 233 | }, |
| 234 | count: row.n, |
| 235 | }) |
| 236 | }) |
| 237 | .collect(); |
| 238 | if readable.is_empty() { |
| 239 | return Ok(Outcome::Ok(Authored::default())); |
| 240 | } |
| 241 | |
| 242 | let all_ids: Vec<&str> = readable.iter().map(|repo| repo.id.as_str()).collect(); |
| 243 | let all = serde_json::to_string(&all_ids)?; |
| 244 | let counts = self |
| 245 | .db |
| 246 | .prepare(COUNTS_SQL) |
| 247 | .bind(&[person.id.as_str().into(), all.as_str().into()])? |
| 248 | .first::<CountsRow>(None) |
| 249 | .await? |
| 250 | .map(|row| AuthoredCounts { |
| 251 | pulls_merged: row.pulls_merged, |
| 252 | pulls_open: row.pulls_open, |
| 253 | pulls: row.pulls, |
| 254 | issues: row.issues, |
| 255 | issues_open: row.issues_open, |
| 256 | }) |
| 257 | .unwrap_or_default(); |
| 258 | |
| 259 | // A repository filter narrows the set; one the viewer cannot read, |
| 260 | // or the person never worked in, leaves nothing. |
| 261 | let shown = match a.repo.as_deref().map(str::trim).filter(|name| !name.is_empty()) { |
| 262 | Some(name) => serde_json::to_string( |
| 263 | &readable |
| 264 | .iter() |
| 265 | .filter(|repo| is_named(repo, name)) |
| 266 | .map(|repo| repo.id.as_str()) |
| 267 | .collect::<Vec<_>>(), |
| 268 | )?, |
| 269 | None => all, |
| 270 | }; |
| 271 | let limit = a.limit.unwrap_or(AUTHORED_PAGE).clamp(1, AUTHORED_PAGE); |
| 272 | let mut params = Params(Vec::new()); |
| 273 | let sql = page_sql(&a, &mut params, &person.id, &shown, limit); |
| 274 | let mut rows = self |
| 275 | .db |
| 276 | .prepare(sql) |
| 277 | .bind(¶ms.values())? |
| 278 | .all() |
| 279 | .await? |
| 280 | .results::<ItemRow>()?; |
| 281 | let next = if rows.len() > limit as usize { |
| 282 | rows.truncate(limit as usize); |
| 283 | rows.last().map(|row| cursor(a.sort, row)) |
| 284 | } else { |
| 285 | None |
| 286 | }; |
| 287 | let items = rows |
| 288 | .into_iter() |
| 289 | .filter_map(|row| { |
| 290 | let repo = by_id.get(row.repo_id.as_str())?; |
| 291 | Some(AuthoredItem { |
| 292 | kind: row.kind, |
| 293 | repo: RepoPath { |
| 294 | namespace: repo.namespace.clone(), |
| 295 | name: repo.name.clone(), |
| 296 | }, |
| 297 | number: row.number, |
| 298 | title: row.title, |
| 299 | state: row.state, |
| 300 | draft: row.status == Some(PullStatus::Draft), |
| 301 | merged: row.status == Some(PullStatus::Merged), |
| 302 | status: row.status, |
| 303 | reason: row.reason, |
| 304 | created_at: row.created_at, |
| 305 | updated_at: row.updated_at, |
| 306 | merged_at: row.merged_at, |
| 307 | }) |
| 308 | }) |
| 309 | .collect(); |
| 310 | Ok(Outcome::Ok(Authored { |
| 311 | items, |
| 312 | next, |
| 313 | counts, |
| 314 | repos, |
| 315 | })) |
| 316 | } |
| 317 | } |
| 318 | |
| 319 | #[cfg(test)] |
| 320 | mod tests { |
| 321 | use super::*; |
| 322 | |
| 323 | fn args() -> ByAuthorArgs { |
| 324 | ByAuthorArgs { |
| 325 | username: "ada".into(), |
| 326 | viewer: None, |
| 327 | kind: None, |
| 328 | state: None, |
| 329 | repo: None, |
| 330 | sort: AuthoredSort::Created, |
| 331 | before: None, |
| 332 | limit: None, |
| 333 | } |
| 334 | } |
| 335 | |
| 336 | fn row(created: &str, updated: &str, id: &str) -> ItemRow { |
| 337 | ItemRow { |
| 338 | kind: AuthoredKind::Pull, |
| 339 | id: id.into(), |
| 340 | repo_id: "rep_1".into(), |
| 341 | number: 1, |
| 342 | title: "t".into(), |
| 343 | state: State::Open, |
| 344 | reason: None, |
| 345 | status: Some(PullStatus::Open), |
| 346 | created_at: created.into(), |
| 347 | updated_at: updated.into(), |
| 348 | merged_at: None, |
| 349 | } |
| 350 | } |
| 351 | |
| 352 | #[test] |
| 353 | fn every_query_is_confined_to_the_readable_set() { |
| 354 | let mut params = Params(Vec::new()); |
| 355 | let sql = page_sql(&args(), &mut params, "usr_1", "[]", 25); |
| 356 | assert_eq!(sql.matches("json_each(?2)").count(), 2); |
| 357 | assert_eq!(sql.matches("author_id = ?1").count(), 2); |
| 358 | assert_eq!(COUNTS_SQL.matches("json_each(?2)").count(), 5); |
| 359 | assert!(sql.contains("LIMIT ?3")); |
| 360 | assert!(sql.contains("ORDER BY created_at DESC, id DESC")); |
| 361 | } |
| 362 | |
| 363 | #[test] |
| 364 | fn filters_by_kind_state_and_page() { |
| 365 | let mut params = Params(Vec::new()); |
| 366 | let a = ByAuthorArgs { |
| 367 | kind: Some(AuthoredKind::Pull), |
| 368 | state: Some(AuthoredState::Merged), |
| 369 | sort: AuthoredSort::Oldest, |
| 370 | before: Some("2026-10-01T00:00:00.000Z|pul_9".into()), |
| 371 | ..args() |
| 372 | }; |
| 373 | let sql = page_sql(&a, &mut params, "usr_1", "[]", 10); |
| 374 | assert!(sql.contains("kind = ?3")); |
| 375 | assert!(sql.contains("status = 'merged'")); |
| 376 | assert!(sql.contains("(created_at > ?4 OR (created_at = ?4 AND id > ?5))")); |
| 377 | assert!(sql.contains("ORDER BY created_at ASC, id ASC")); |
| 378 | assert!(sql.contains("LIMIT ?6")); |
| 379 | assert_eq!(params.0.len(), 6); |
| 380 | } |
| 381 | |
| 382 | #[test] |
| 383 | fn ignores_a_cursor_that_is_not_one() { |
| 384 | let mut params = Params(Vec::new()); |
| 385 | let a = ByAuthorArgs { |
| 386 | before: Some("nonsense".into()), |
| 387 | ..args() |
| 388 | }; |
| 389 | let sql = page_sql(&a, &mut params, "usr_1", "[]", 10); |
| 390 | assert!(!sql.contains(" OR (")); |
| 391 | assert_eq!(parse_cursor("a|b|c"), None); |
| 392 | assert_eq!(parse_cursor("|b"), None); |
| 393 | assert_eq!(parse_cursor("k|i"), Some(("k", "i"))); |
| 394 | } |
| 395 | |
| 396 | #[test] |
| 397 | fn a_cursor_names_the_sort_key() { |
| 398 | let r = row("2026-01-01T00:00:00.000Z", "2026-02-01T00:00:00.000Z", "pul_1"); |
| 399 | assert_eq!(cursor(AuthoredSort::Created, &r), "2026-01-01T00:00:00.000Z|pul_1"); |
| 400 | assert_eq!(cursor(AuthoredSort::Updated, &r), "2026-02-01T00:00:00.000Z|pul_1"); |
| 401 | assert_eq!(parse_cursor(&cursor(AuthoredSort::Oldest, &r)), Some(("2026-01-01T00:00:00.000Z", "pul_1"))); |
| 402 | } |
| 403 | |
| 404 | #[test] |
| 405 | fn names_a_repository_by_its_path() { |
| 406 | let repo = Repo { |
| 407 | id: "rep_1".into(), |
| 408 | namespace: "acme".into(), |
| 409 | name: "Rocket".into(), |
| 410 | description: None, |
| 411 | is_private: false, |
| 412 | owner_id: "usr_1".into(), |
| 413 | default_branch: "main".into(), |
| 414 | fork_of: None, |
| 415 | protected: false, |
| 416 | created_at: String::new(), |
| 417 | topics: Vec::new(), |
| 418 | website: None, |
| 419 | archived_at: None, |
| 420 | }; |
| 421 | assert!(is_named(&repo, "acme/rocket")); |
| 422 | assert!(is_named(&repo, "ACME/Rocket")); |
| 423 | assert!(!is_named(&repo, "acme")); |
| 424 | assert!(!is_named(&repo, "other/rocket")); |
| 425 | } |
| 426 | } |