Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 1 | //! Reading a page's worth of rows in one round trip. |
| 2 | //! | |
| 3 | //! A pull request's page used to cost about twenty queries one after | |
| 4 | //! another: the repository, the pull request, its issue, then where it | |
| 5 | //! stands (its progress, latest review, settings, statuses, verdicts, | |
| 6 | //! queue entry, the confidence signals), its comments, checks and | |
| 7 | //! messages. Each is a round trip to the database, and from most places a | |
| 8 | //! round trip is 20 to 40 ms. Here they go as one D1 batch, keyed by the | |
| 9 | //! repository's id and the pull request's number, and the batch starts | |
| 10 | //! while the repos service is still deciding whether the viewer may see | |
| 11 | //! the repository ([`Work::repo_then`]). | |
| 12 | //! | |
| 13 | //! The helpers that answer those questions (`settings`, `statuses`, | |
| 14 | //! `review_pending`, `approvals_gap` and the rest) look here first, so the | |
| 15 | //! code that decides a pull request's lifecycle is the same code whether | |
| 16 | //! its rows were batched or read one at a time. Only reads are kept; a | |
| 17 | //! helper that writes and reads back (mergeability) reads the database. | |
| 18 | ||
| 19 | use std::cell::RefCell; | |
| 20 | use std::collections::HashMap; | |
| 21 | use std::future::Future; | |
| 22 | use std::rc::Rc; | |
| 23 | ||
| 24 | use g1t_contracts::repos::{Repo, RepoPath}; | |
| 25 | use g1t_contracts::time::rfc3339; | |
| 26 | use g1t_contracts::{Outcome, Viewer}; | |
| 27 | use g1t_kit::now_ms; | |
| 28 | use serde::de::DeserializeOwned; | |
| 29 | use worker::wasm_bindgen::JsValue; | |
| 30 | use worker::{D1PreparedStatement, D1Result, Result}; | |
| 31 | ||
| 32 | use crate::reviews::AGENT_ID; | |
| 33 | use crate::rows::PULL_COLUMNS; | |
| 34 | use crate::{ISSUE_COLUMNS, Work}; | |
| 35 | ||
| 36 | /// The pull request's id, from the batch's `?1` (repository) and `?2` (number). | |
| 37 | const PULL_ID: &str = "(SELECT id FROM pulls WHERE repo_id = ?1 AND number = ?2)"; | |
| 38 | ||
| 39 | /// How long a review that never reported is waited for (reviews.rs). | |
| 40 | pub(crate) const REVIEW_PENDING_MS: u64 = 30 * 60 * 1000; | |
| 41 | ||
| 42 | /// Each statement of a pull request's batch, in order. | |
| 43 | #[derive(Clone, Copy)] | |
| 44 | pub(crate) enum Slot { | |
| 45 | Pull, | |
| 46 | Issue, | |
| 47 | Comments, | |
| 48 | Verdicts, | |
| 49 | Runs, | |
| 50 | RunHistory, | |
| 51 | Review, | |
| 52 | ReviewComments, | |
| 53 | ReviewPending, | |
| 54 | Statuses, | |
| 55 | Settings, | |
| 56 | Hold, | |
| 57 | Messages, | |
| 58 | Others, | |
| 59 | Queued, | |
| 60 | LatestRun, | |
| 61 | Halted, | |
| 62 | Denials, | |
| 63 | Unanswered, | |
| 64 | Planned, | |
| 65 | } | |
| 66 | ||
| 67 | /// How many runs `latest_checks` and `earlier_checks` show together. | |
| 68 | pub(crate) const RECENT_RUNS: u32 = 11; | |
| 69 | ||
| 70 | /// One pull request's rows, as read for this request. | |
| 71 | pub(crate) struct Prefetched { | |
| 72 | pub(crate) pull_id: String, | |
| 73 | pub(crate) repo_id: String, | |
| 74 | /// Its head commit when read, which its statuses are for. | |
| 75 | pub(crate) head: Option<String>, | |
| 76 | results: Vec<D1Result>, | |
| 77 | } | |
| 78 | ||
| 79 | impl Prefetched { | |
| 80 | /// A slot's rows, as the helper that reads them would have. | |
| 81 | pub(crate) fn rows<T: DeserializeOwned>(&self, slot: Slot) -> Result<Vec<T>> { | |
| 82 | match self.results.get(slot as usize) { | |
| 83 | Some(result) => result.results::<T>(), | |
| 84 | None => Ok(Vec::new()), | |
| 85 | } | |
| 86 | } | |
| 87 | ||
| 88 | /// A slot's first row. | |
| 89 | pub(crate) fn first<T: DeserializeOwned>(&self, slot: Slot) -> Result<Option<T>> { | |
| 90 | Ok(self.rows::<T>(slot)?.into_iter().next()) | |
| 91 | } | |
| 92 | } | |
| 93 | ||
| 94 | #[derive(serde::Deserialize)] | |
| 95 | struct IdRow { | |
| 96 | id: String, | |
| 97 | head_commit: Option<String>, | |
| 98 | } | |
| 99 | ||
| 100 | thread_local! { | |
| 101 | /// Repository ids by path, as the repos service last answered: a guess | |
| 102 | /// that lets a page's batch start before the answer comes. Never | |
| 103 | /// trusted: nothing read with a guess is used unless the repos service | |
| 104 | /// then says the viewer may see that very repository. | |
| 105 | static REPO_IDS: RefCell<HashMap<String, String>> = RefCell::new(HashMap::new()); | |
| 106 | } | |
| 107 | ||
| 108 | /// Guesses kept at most; the map is emptied past this. | |
| 109 | const MAX_GUESSES: usize = 1024; | |
| 110 | ||
| 111 | fn path_key(path: &RepoPath) -> String { | |
| 112 | format!("{}/{}", path.namespace.to_lowercase(), path.name.to_lowercase()) | |
| 113 | } | |
| 114 | ||
| 115 | fn guess(path: &RepoPath) -> Option<String> { | |
| 116 | REPO_IDS.with(|ids| ids.borrow().get(&path_key(path)).cloned()) | |
| 117 | } | |
| 118 | ||
| 119 | fn learn(path: &RepoPath, id: Option<&str>) { | |
| 120 | REPO_IDS.with(|ids| { | |
| 121 | let mut ids = ids.borrow_mut(); | |
| 122 | match id { | |
| 123 | Some(id) => { | |
| 124 | if ids.len() >= MAX_GUESSES { | |
| 125 | ids.clear(); | |
| 126 | } | |
| 127 | ids.insert(path_key(path), id.to_owned()); | |
| 128 | } | |
| 129 | None => { | |
| 130 | ids.remove(&path_key(path)); | |
| 131 | } | |
| 132 | } | |
| 133 | }); | |
| 134 | } | |
| 135 | ||
| 136 | impl Work { | |
| 137 | /// The repository at `path` if `viewer` may see it, and what `load` | |
| 138 | /// reads for it. When this isolate has seen the repository before, | |
| 139 | /// `load` starts with its id at once, beside the access check, instead | |
| 140 | /// of after it: one round trip instead of two. What it read is | |
| 141 | /// discarded unless the check passes for that same repository. | |
| 142 | pub(crate) async fn repo_then<T, F, Fut>(&self, path: &RepoPath, viewer: &Viewer, load: F) -> Result<Outcome<(Repo, T)>> | |
| 143 | where | |
| 144 | F: Fn(String) -> Fut, | |
| 145 | Fut: Future<Output = Result<T>>, | |
| 146 | { | |
| 147 | let guessed = guess(path); | |
| 148 | let (found, early) = match &guessed { | |
| 149 | Some(id) => { | |
| 150 | let (found, early) = futures_util::future::join(self.repo(path, viewer), load(id.clone())).await; | |
| 151 | (found?, Some(early)) | |
| 152 | } | |
| 153 | None => (self.repo(path, viewer).await?, None), | |
| 154 | }; | |
| 155 | let repo = match found { | |
| 156 | Outcome::Ok(repo) => repo, | |
| 157 | Outcome::Fail(failure) => { | |
| 158 | learn(path, None); | |
| 159 | return Ok(Outcome::Fail(failure)); | |
| 160 | } | |
| 161 | }; | |
| 162 | if let (Some(id), Some(Ok(value))) = (&guessed, early) | |
| 163 | && *id == repo.id | |
| 164 | { | |
| 165 | return Ok(Outcome::Ok((repo, value))); | |
| 166 | } | |
| 167 | learn(path, Some(&repo.id)); | |
| 168 | let value = load(repo.id.clone()).await?; | |
| 169 | Ok(Outcome::Ok((repo, value))) | |
| 170 | } | |
| 171 | ||
| 172 | fn statement(&self, sql: &str, binds: &[JsValue]) -> Result<D1PreparedStatement> { | |
| 173 | self.db.prepare(sql).bind(binds) | |
| 174 | } | |
| 175 | ||
| 176 | /// Everything a pull request's page and its lifecycle read, in one | |
| 177 | /// batch. `None` when there is no such pull request. | |
| 178 | pub(crate) async fn prefetch_pull(&self, repo_id: String, number: u32) -> Result<Option<Prefetched>> { | |
| 179 | let key = || -> [JsValue; 2] { [repo_id.as_str().into(), number.into()] }; | |
| 180 | let with = |extra: JsValue| -> [JsValue; 3] { [repo_id.as_str().into(), number.into(), extra] }; | |
| 181 | let repo_only = || -> [JsValue; 1] { [repo_id.as_str().into()] }; | |
| 182 | let statements = vec![ | |
| 183 | // Slot::Pull: the row, with everything the lifecycle tracks on it. | |
| 184 | self.statement(&format!("SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ?1 AND number = ?2"), &key())?, | |
| 185 | // Slot::Issue | |
| 186 | self.statement( | |
| 187 | &format!( | |
| 188 | "SELECT {ISSUE_COLUMNS} FROM issues | |
| 189 | WHERE repo_id = ?1 AND number = (SELECT issue_number FROM pulls WHERE repo_id = ?1 AND number = ?2)" | |
| 190 | ), | |
| 191 | &key(), | |
| 192 | )?, | |
| 193 | // Slot::Comments (lib.rs `comments`) | |
| 194 | self.statement("SELECT * FROM comments WHERE repo_id = ?1 AND number = ?2 ORDER BY id LIMIT 500", &key())?, | |
| 195 | // Slot::Verdicts: every verdict, oldest first, for the approval | |
| 196 | // rule, a person's request for changes and a person's approval. | |
| 197 | self.statement( | |
| 198 | "SELECT author_id, author_name, verdict, created_at FROM comments | |
| 199 | WHERE repo_id = ?1 AND number = ?2 AND verdict IS NOT NULL ORDER BY id", | |
| 200 | &key(), | |
| 201 | )?, | |
| 202 | // Slot::Runs: the latest and the ten before it (checks.rs). | |
| 203 | self.statement( | |
| 204 | &format!("SELECT * FROM check_runs WHERE pull_id = {PULL_ID} ORDER BY id DESC LIMIT ?3"), | |
| 205 | &with(RECENT_RUNS.into()), | |
| 206 | )?, | |
| 207 | // Slot::RunHistory (confidence.rs `signals`) | |
| 208 | self.statement( | |
| 209 | &format!("SELECT head_commit, status FROM check_runs WHERE pull_id = {PULL_ID} ORDER BY id LIMIT 50"), | |
| 210 | &key(), | |
| 211 | )?, | |
| 212 | // Slot::Review (lifecycle.rs `assess_now`) | |
| 213 | self.statement( | |
| 214 | &format!( | |
| 215 | "SELECT finished_at, verdict FROM review_runs | |
| 216 | WHERE pull_id = {PULL_ID} AND finished_at IS NOT NULL ORDER BY id DESC LIMIT 1" | |
| 217 | ), | |
| 218 | &key(), | |
| 219 | )?, | |
| 220 | // Slot::ReviewComments: lines that latest review commented on | |
| 221 | // (confidence.rs `review_comments`). | |
| 222 | self.statement( | |
| 223 | &format!( | |
| 224 | "SELECT count(*) AS n FROM comments | |
| 225 | WHERE repo_id = ?1 AND number = ?2 AND author_id = ?3 AND path IS NOT NULL | |
| 226 | AND created_at = (SELECT finished_at FROM review_runs | |
| 227 | WHERE pull_id = {PULL_ID} AND finished_at IS NOT NULL ORDER BY id DESC LIMIT 1)" | |
| 228 | ), | |
| 229 | &with(AGENT_ID.into()), | |
| 230 | )?, | |
| 231 | // Slot::ReviewPending (reviews.rs `review_pending`) | |
| 232 | self.statement( | |
| 233 | &format!( | |
| 234 | "SELECT id AS value FROM review_runs | |
| 235 | WHERE pull_id = {PULL_ID} AND finished_at IS NULL AND created_at > ?3 LIMIT 1" | |
| 236 | ), | |
| 237 | &with(rfc3339(now_ms().saturating_sub(REVIEW_PENDING_MS)).into()), | |
| 238 | )?, | |
| 239 | // Slot::Statuses: on its head (statuses.rs `statuses`). | |
| 240 | self.statement( | |
| 241 | "SELECT context, state, description, target_url, updated_at FROM commit_statuses | |
| 242 | WHERE repo_id = ?1 AND sha = (SELECT head_commit FROM pulls WHERE repo_id = ?1 AND number = ?2) | |
| 243 | ORDER BY context", | |
| 244 | &key(), | |
| 245 | )?, | |
| 246 | // Slot::Settings and Slot::Hold (settings.rs `settings`) | |
| 247 | self.statement("SELECT * FROM repo_settings WHERE repo_id = ?1", &repo_only())?, | |
| 248 | self.statement("SELECT hold_low AS n FROM confidence_rules WHERE repo_id = ?1", &repo_only())?, | |
| 249 | // Slot::Messages (messages.rs `messages`) | |
| 250 | self.statement( | |
| 251 | &format!("SELECT * FROM agent_messages WHERE pull_id = {PULL_ID} ORDER BY created_at, id"), | |
| 252 | &key(), | |
| 253 | )?, | |
| 254 | // Slot::Others (reviews.rs `overlaps`) | |
| 255 | self.statement( | |
| 256 | "SELECT number, title, issue_number, files FROM pulls | |
| 257 | WHERE repo_id = ?1 AND number != ?2 AND status IN ('draft', 'open') | |
| 258 | ORDER BY number LIMIT 200", | |
| 259 | &key(), | |
| 260 | )?, | |
| 261 | // Slot::Queued (queue.rs `queued_entry`) | |
| 262 | self.statement( | |
| 263 | &format!( | |
| 264 | "SELECT * FROM queue_entries | |
| 265 | WHERE pull_id = {PULL_ID} AND state IN ('waiting', 'testing', 'passed') LIMIT 1" | |
| 266 | ), | |
| 267 | &key(), | |
| 268 | )?, | |
| 269 | // Slot::LatestRun, Halted, Denials, Unanswered, Planned | |
| 270 | // (confidence.rs `signals`) | |
| 271 | self.statement( | |
| 272 | &format!( | |
| 273 | "SELECT r.id, r.cost_usd, r.budget_usd, r.time_cap_minutes, | |
| 274 | (julianday(COALESCE(r.finished_at, r.updated_at)) - julianday(COALESCE(r.started_at, r.created_at))) * 1440 AS minutes, | |
| 275 | c.self_level, c.uncertain_about | |
| 276 | FROM agent_runs r LEFT JOIN run_confidence c ON c.run_id = r.id | |
| 277 | WHERE r.pull_id = {PULL_ID} AND r.kind IN ('implement', 'revise') | |
| 278 | ORDER BY r.created_at DESC LIMIT 1" | |
| 279 | ), | |
| 280 | &key(), | |
| 281 | )?, | |
| 282 | self.statement( | |
| 283 | &format!( | |
| 284 | "SELECT halted FROM agent_runs WHERE pull_id = {PULL_ID} AND halted IS NOT NULL | |
| 285 | ORDER BY created_at DESC LIMIT 1" | |
| 286 | ), | |
| 287 | &key(), | |
| 288 | )?, | |
| 289 | self.statement( | |
| 290 | &format!( | |
| 291 | "SELECT count(*) AS n FROM session_entries | |
| 292 | WHERE pull_id = {PULL_ID} AND kind = 'note' AND text LIKE 'Denied:%'" | |
| 293 | ), | |
| 294 | &key(), | |
| 295 | )?, | |
| 296 | self.statement( | |
| 297 | "SELECT count(*) AS n FROM agent_messages | |
| 298 | WHERE repo_id = ?1 AND from_number = ?2 AND kind IN ('question', 'handoff') | |
| 299 | AND answered_at IS NULL", | |
| 300 | &key(), | |
| 301 | )?, | |
| 302 | self.statement( | |
| 303 | "SELECT json_extract(planned.value, '$.files') AS files | |
| 304 | FROM plans, json_each(plans.issues) AS planned | |
| 305 | WHERE plans.repo_id = ?1 AND plans.status = 'applied' | |
| 306 | AND json_extract(planned.value, '$.number') = (SELECT issue_number FROM pulls WHERE repo_id = ?1 AND number = ?2) | |
| 307 | LIMIT 1", | |
| 308 | &key(), | |
| 309 | )?, | |
| 310 | ]; | |
| 311 | let count = statements.len() as u32; | |
| 312 | let results = self.timing.db(count, self.db.batch(statements)).await?; | |
| 313 | let Some(row) = results | |
| 314 | .first() | |
| 315 | .map(|result| result.results::<IdRow>()) | |
| 316 | .transpose()? | |
| 317 | .and_then(|rows| rows.into_iter().next()) | |
| 318 | else { | |
| 319 | return Ok(None); | |
| 320 | }; | |
| 321 | Ok(Some(Prefetched { | |
| 322 | pull_id: row.id, | |
| 323 | repo_id, | |
| 324 | head: row.head_commit, | |
| 325 | results, | |
| 326 | })) | |
| 327 | } | |
| 328 | ||
| 329 | /// The rows read for pull request `pull_id` in this request, if any. | |
| 330 | pub(crate) fn prefetched_pull(&self, pull_id: &str) -> Option<Rc<Prefetched>> { | |
| 331 | self.prefetched | |
| 332 | .borrow() | |
| 333 | .as_ref() | |
| 334 | .filter(|found| found.pull_id == pull_id) | |
| 335 | .cloned() | |
| 336 | } | |
| 337 | ||
| 338 | /// The rows read in this request for a pull request of `repo_id`. | |
| 339 | pub(crate) fn prefetched_repo(&self, repo_id: &str) -> Option<Rc<Prefetched>> { | |
| 340 | self.prefetched | |
| 341 | .borrow() | |
| 342 | .as_ref() | |
| 343 | .filter(|found| found.repo_id == repo_id) | |
| 344 | .cloned() | |
| 345 | } | |
| 346 | ||
| 347 | /// Keeps `found` for the rest of this request. | |
| 348 | pub(crate) fn keep_prefetched(&self, found: Option<Prefetched>) { | |
| 349 | *self.prefetched.borrow_mut() = found.map(Rc::new); | |
| 350 | } | |
| 351 | } | |
| 352 | ||
| 353 | #[cfg(test)] | |
| 354 | mod tests { | |
| 355 | use super::*; | |
| 356 | ||
| 357 | fn at(namespace: &str, name: &str) -> RepoPath { | |
| 358 | RepoPath { namespace: namespace.to_owned(), name: name.to_owned() } | |
| 359 | } | |
| 360 | ||
| 361 | #[test] | |
| 362 | fn guesses_ignore_case_and_are_forgotten_on_refusal() { | |
| 363 | learn(&at("Flagon-IO", "G1T"), Some("repo_1")); | |
| 364 | assert_eq!(guess(&at("flagon-io", "g1t")).as_deref(), Some("repo_1")); | |
| 365 | learn(&at("flagon-io", "g1t"), None); | |
| 366 | assert_eq!(guess(&at("flagon-io", "g1t")), None); | |
| 367 | } | |
| 368 | } |