From e1caeaf070e3ee55887d700b59d9da99e24a9f57 Mon Sep 17 00:00:00 2001 From: "Jonathan D.A. Jewell" <6759885+hyperpolymath@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:23:08 +0100 Subject: [PATCH] feat: `squabble board` and `squabble inbox-sweep` MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `board` enumerates every non-archived repo of each owner, reads the merge gate (rulesets ∪ classic protection) and open PRs of repos that have any, and places each PR in one bucket: needs you / agent work / landing by itself / stale. Renders Markdown under the 65536-byte issue body limit and, with --publish owner/repo#N, rewrites that issue's body (never comments) and checks the stored body equals the one sent. A whole-batch failure (502, no data) is split and retried down to one repo; listing pages retry transport failures. Exit 6 = incomplete, with the unreadable repos named at the top of the board. `inbox-sweep` lists every notification thread still in the inbox, resolves PR/issue state in GraphQL batches, and clears (unsubscribe, then mark done) only threads whose subject is merged or closed. Dry run by default; --apply writes, stops at a REST reserve, and names every due thread it left. Unknown state and non-PR/issue subjects are kept. Live board 2026-10-01: 398+50 repos, 458+54 open PRs, all placed, 0 unreadable; counts reconcile with `gh search prs` (the one remaining difference is a PR in an archived repo). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_0196aKTfvS6vPQVp2LPoXbjP --- crates/squabble-cli/src/board.rs | 407 ++++++++++ crates/squabble-cli/src/inbox.rs | 542 +++++++++++++ crates/squabble-cli/src/main.rs | 15 +- crates/squabble-core/src/board.rs | 626 +++++++++++++++ crates/squabble-core/src/inbox.rs | 225 ++++++ crates/squabble-core/src/lib.rs | 2 + .../squabble-forge/graphql/board_repo.graphql | 53 ++ .../graphql/estate_repos.graphql | 21 + crates/squabble-forge/src/board.rs | 730 ++++++++++++++++++ crates/squabble-forge/src/inbox.rs | 209 +++++ crates/squabble-forge/src/lib.rs | 3 + 11 files changed, 2831 insertions(+), 2 deletions(-) create mode 100644 crates/squabble-cli/src/board.rs create mode 100644 crates/squabble-cli/src/inbox.rs create mode 100644 crates/squabble-core/src/board.rs create mode 100644 crates/squabble-core/src/inbox.rs create mode 100644 crates/squabble-forge/graphql/board_repo.graphql create mode 100644 crates/squabble-forge/graphql/estate_repos.graphql create mode 100644 crates/squabble-forge/src/board.rs create mode 100644 crates/squabble-forge/src/inbox.rs diff --git a/crates/squabble-cli/src/board.rs b/crates/squabble-cli/src/board.rs new file mode 100644 index 0000000..43ad9c1 --- /dev/null +++ b/crates/squabble-cli/src/board.rs @@ -0,0 +1,407 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +//! `squabble board` — the estate "needs me" board. +//! +//! Enumerates every non-archived repository of each owner, reads the merge +//! gate and open PRs of the repos that have any, places each PR in one bucket +//! (`squabble_core::board::classify`) and renders Markdown — or JSON with +//! `--json`. +//! +//! `--publish owner/repo#N` rewrites the *body* of issue N with the board. It +//! never comments, so a run never notifies anyone. That edit is the only write +//! this subcommand can make. +//! +//! Exit `0` = complete board; `6` = board produced but incomplete (a repo was +//! unreadable or had more open PRs than one page) — it is still rendered and +//! still published, with the gap stated at the top; `2` = usage or tool +//! failure, including any owner whose enumeration failed. + +use squabble_core::board::{ + classify, epoch_day, render_markdown, rfc3339_from_unix, Board, ISSUE_BODY_LIMIT, +}; +use squabble_core::chains::RepoId; +use squabble_forge::board::{enumerate_owner, fetch_board, RepoRead, BOARD_BATCH}; +use squabble_forge::{GhTransport, GraphQlTransport}; +use std::io::Write; +use std::process::{Command, ExitCode, Stdio}; +use std::time::{SystemTime, UNIX_EPOCH}; + +pub(crate) const USAGE: &str = "usage: squabble board [--owners a,b] [--stale-days N] [--batch N] [--json] [--out FILE] [--publish owner/repo#N]"; + +/// Exit code for a board that was produced but has stated gaps. +pub(crate) const INCOMPLETE_EXIT: u8 = 6; + +/// Owners read when `--owners` is not given. +const DEFAULT_OWNERS: &[&str] = &["hyperpolymath", "metadatastician"]; + +struct Args { + owners: Vec, + stale_days: i64, + batch: usize, + json: bool, + out: Option, + publish: Option<(String, u64)>, +} + +/// Parse `owner/repo#N`. +fn parse_issue_ref(s: &str) -> Option<(String, u64)> { + let (slug, n) = s.split_once('#')?; + let (o, r) = slug.split_once('/')?; + if o.is_empty() || r.is_empty() || r.contains('/') { + return None; + } + Some((slug.to_string(), n.parse().ok().filter(|&n: &u64| n > 0)?)) +} + +/// Parse the flags after `squabble board`. +fn parse_args(rest: &[String]) -> Result { + let mut a = Args { + owners: DEFAULT_OWNERS.iter().map(|s| s.to_string()).collect(), + stale_days: 30, + batch: BOARD_BATCH, + json: false, + out: None, + publish: None, + }; + let mut it = rest.iter(); + while let Some(arg) = it.next() { + match arg.as_str() { + "--json" => a.json = true, + "--owners" => { + a.owners = it + .next() + .map(|v| { + v.split(',') + .map(str::trim) + .filter(|s| !s.is_empty()) + .map(str::to_string) + .collect::>() + }) + .filter(|v| !v.is_empty()) + .ok_or("--owners needs a comma-separated list")? + } + "--stale-days" => { + a.stale_days = it + .next() + .and_then(|v| v.parse().ok()) + .filter(|&n: &i64| n > 0) + .ok_or("--stale-days needs a positive integer")? + } + "--batch" => { + a.batch = it + .next() + .and_then(|v| v.parse().ok()) + .filter(|&n: &usize| n > 0) + .ok_or("--batch needs a positive integer")? + } + "--out" => a.out = Some(it.next().ok_or("--out needs a file path")?.clone()), + "--publish" => { + a.publish = Some( + it.next() + .and_then(|v| parse_issue_ref(v)) + .ok_or("--publish needs owner/repo#N")?, + ) + } + s => return Err(format!("unexpected argument `{s}`\n{USAGE}")), + } + } + Ok(a) +} + +/// Build the board for `owners` through `transport`. `Err` when any owner +/// cannot be enumerated: a board missing a whole owner would look complete. +fn build( + transport: &dyn GraphQlTransport, + owners: &[String], + batch: usize, + now_secs: u64, + stale_days: i64, +) -> Result { + let generated_at = rfc3339_from_unix(now_secs); + let today = epoch_day(&generated_at).ok_or("clock produced an unparseable date")?; + let mut board = Board { + generated_at, + owners: owners.to_vec(), + ..Board::default() + }; + let mut with_prs: Vec = Vec::new(); + for owner in owners { + let listing = enumerate_owner(transport, owner)?; + board + .repos_enumerated + .insert(owner.clone(), listing.repos.len()); + board + .open_prs_reported + .insert(owner.clone(), listing.open_prs()); + with_prs.extend( + listing + .repos + .into_iter() + .filter(|(_, n)| *n > 0) + .map(|(r, _)| r), + ); + } + let (reads, _queries) = fetch_board(transport, &with_prs, batch); + for (repo, read) in with_prs.iter().zip(reads) { + match read { + RepoRead::Unavailable(why) => board.unavailable.push((repo.clone(), why)), + RepoRead::Read { + gate, + prs, + truncated, + } => { + if truncated { + board.truncated.push(repo.clone()); + } + for pr in prs { + let reason = classify(&pr, &gate, today, stale_days); + board.placed.push((pr, reason)); + } + } + } + } + Ok(board) +} + +/// Replace the body of `slug#number` with `body` via `gh api`, then confirm +/// the stored body is the one sent. +fn publish(slug: &str, number: u64, body: &str) -> Result<(), String> { + let payload = serde_json::json!({ "body": body }).to_string(); + let mut child = Command::new("gh") + .args([ + "api", + "-X", + "PATCH", + &format!("repos/{slug}/issues/{number}"), + "--input", + "-", + ]) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .map_err(|e| format!("failed to run `gh api`: {e}"))?; + child + .stdin + .as_mut() + .ok_or("could not open stdin for `gh`")? + .write_all(payload.as_bytes()) + .map_err(|e| format!("could not write to `gh`: {e}"))?; + let out = child + .wait_with_output() + .map_err(|e| format!("`gh api` did not finish: {e}"))?; + if !out.status.success() { + return Err(format!( + "PATCH {slug}#{number} failed ({}): {}", + out.status, + String::from_utf8_lossy(&out.stderr).trim() + )); + } + // rc=0 is not evidence: check what GitHub actually stored. + let stored: serde_json::Value = serde_json::from_slice(&out.stdout) + .map_err(|e| format!("PATCH {slug}#{number}: response was not JSON: {e}"))?; + match stored.get("body").and_then(serde_json::Value::as_str) { + Some(b) if b == body => Ok(()), + Some(b) => Err(format!( + "PATCH {slug}#{number}: stored body differs from the one sent ({} vs {} bytes)", + b.len(), + body.len() + )), + None => Err(format!("PATCH {slug}#{number}: response carried no body")), + } +} + +/// Entry point for `squabble board `. +pub(crate) fn run(rest: &[String]) -> ExitCode { + let args = match parse_args(rest) { + Ok(a) => a, + Err(e) => { + eprintln!("squabble board: {e}"); + return ExitCode::from(2); + } + }; + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0); + let board = match build(&GhTransport, &args.owners, args.batch, now, args.stale_days) { + Ok(b) => b, + Err(e) => { + eprintln!("squabble board: {e}"); + return ExitCode::from(2); + } + }; + let md = render_markdown(&board, ISSUE_BODY_LIMIT); + let text = if args.json { + match serde_json::to_string_pretty(&board) { + Ok(j) => j, + Err(e) => { + eprintln!("squabble board: could not serialise: {e}"); + return ExitCode::from(2); + } + } + } else { + md.clone() + }; + match &args.out { + Some(path) => { + if let Err(e) = std::fs::write(path, &text) { + eprintln!("squabble board: cannot write `{path}`: {e}"); + return ExitCode::from(2); + } + } + None => println!("{text}"), + } + let placed = board.placed_per_owner(); + for o in &board.owners { + eprintln!( + "squabble board: {o}: {} repos, {} open PRs reported, {} placed", + board.repos_enumerated.get(o).copied().unwrap_or(0), + board.open_prs_reported.get(o).copied().unwrap_or(0), + placed.get(o).copied().unwrap_or(0) + ); + } + if let Some((slug, n)) = &args.publish { + if let Err(e) = publish(slug, *n, &md) { + eprintln!("squabble board: {e}"); + return ExitCode::from(2); + } + eprintln!( + "squabble board: published to {slug}#{n} ({} bytes)", + md.len() + ); + } + if board.unavailable.is_empty() && board.truncated.is_empty() { + ExitCode::SUCCESS + } else { + eprintln!( + "squabble board: INCOMPLETE — {} unreadable, {} truncated", + board.unavailable.len(), + board.truncated.len() + ); + ExitCode::from(INCOMPLETE_EXIT) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::{json, Value}; + use squabble_core::board::{Bucket, Reason}; + use std::cell::RefCell; + + fn args(v: &[&str]) -> Vec { + v.iter().map(|s| s.to_string()).collect() + } + + #[test] + fn publish_ref_parses_and_rejects() { + assert_eq!( + parse_issue_ref("hyperpolymath/standards#12"), + Some(("hyperpolymath/standards".into(), 12)) + ); + for bad in ["standards#12", "a/b#0", "a/b#x", "a/b/c#1", "/b#1", "a/b"] { + assert_eq!(parse_issue_ref(bad), None, "{bad}"); + } + } + + #[test] + fn flags_parse_and_unknown_flags_fail() { + let a = parse_args(&args(&["--owners", "x, y", "--stale-days", "7", "--json"])).unwrap(); + assert_eq!(a.owners, vec!["x", "y"]); + assert_eq!(a.stale_days, 7); + assert!(a.json); + assert!(parse_args(&args(&["--nope"])).is_err()); + assert!(parse_args(&args(&["--publish", "bad"])).is_err()); + } + + /// Answers listing queries and board queries from fixed fixtures. + struct Fake { + calls: RefCell, + } + + impl GraphQlTransport for Fake { + fn execute(&self, body: &Value) -> Result { + *self.calls.borrow_mut() += 1; + let rate = json!({ "cost": 1, "remaining": 4000, "resetAt": null }); + if body["variables"].get("login").is_some() { + return Ok( + json!({ "data": { "rateLimit": rate, "repositoryOwner": { "repositories": { + "totalCount": 3, + "pageInfo": { "hasNextPage": false, "endCursor": null }, + "nodes": [ + { "nameWithOwner": "me/gated", "pullRequests": { "totalCount": 2 } }, + { "nameWithOwner": "me/bare", "pullRequests": { "totalCount": 1 } }, + { "nameWithOwner": "me/quiet", "pullRequests": { "totalCount": 0 } } + ] + }}}}), + ); + } + let pr = |n: u64, state: &str, armed: bool| { + json!({ + "number": n, "title": "t", "url": format!("u{n}"), "isDraft": false, + "mergeable": "MERGEABLE", "mergeStateStatus": state, "reviewDecision": null, + "updatedAt": "2026-10-01T00:00:00Z", "headRefName": "h", "author": { "login": "a" }, + "autoMergeRequest": if armed { json!({ "enabledAt": "x" }) } else { Value::Null }, + "commits": { "nodes": [ { "commit": { "statusCheckRollup": { "state": "PENDING" } } } ] } + }) + }; + let mut data = serde_json::Map::new(); + data.insert("rateLimit".into(), rate); + for (i, name) in ["o0", "o1"].iter().enumerate() { + if body["variables"].get(*name).is_none() { + continue; + } + let repo = body["variables"][format!("n{i}")].as_str().unwrap(); + let node = if repo == "gated" { + json!({ "nameWithOwner": "me/gated", + "defaultBranchRef": { "name": "main", "branchProtectionRule": null, + "rules": { "nodes": [ { "type": "REQUIRED_STATUS_CHECKS", + "parameters": { "requiredStatusChecks": [ { "context": "ci" } ] } } ] } }, + "pullRequests": { "totalCount": 2, "nodes": [ pr(1, "CLEAN", false), pr(2, "BLOCKED", true) ] } }) + } else { + json!({ "nameWithOwner": "me/bare", + "defaultBranchRef": { "name": "main", "branchProtectionRule": null, "rules": { "nodes": [] } }, + "pullRequests": { "totalCount": 1, "nodes": [ pr(3, "BLOCKED", false) ] } }) + }; + data.insert(format!("r{i}"), node); + } + Ok(json!({ "data": data })) + } + } + + #[test] + fn planted_controls_land_in_their_buckets_and_counts_reconcile() { + // 2026-10-01T12:00:00Z + let board = build( + &Fake { + calls: RefCell::new(0), + }, + &["me".into()], + 10, + 1_790_856_000, + 30, + ) + .unwrap(); + let reason_of = |n: u64| board.placed.iter().find(|(p, _)| p.number == n).unwrap().1; + assert_eq!(reason_of(1), Reason::CleanMergeByHand); + assert_eq!(reason_of(2), Reason::ArmedWaiting); + assert_eq!(reason_of(3), Reason::NoRequiredGate); + assert_eq!(reason_of(3).bucket(), Bucket::NeedsYou); + // Enumerated vs placed: the set difference is empty. + assert_eq!(board.open_prs_reported["me"], 3); + assert_eq!(board.placed_per_owner()["me"], 3); + assert!(board.unavailable.is_empty() && board.truncated.is_empty()); + assert_eq!(board.repos_enumerated["me"], 3); + } + + #[test] + fn repos_without_open_prs_cost_no_detail_query() { + let fake = Fake { + calls: RefCell::new(0), + }; + build(&fake, &["me".into()], 10, 1_790_856_000, 30).unwrap(); + // One listing page + one board batch (gated, bare); `quiet` not read. + assert_eq!(*fake.calls.borrow(), 2); + } +} diff --git a/crates/squabble-cli/src/inbox.rs b/crates/squabble-cli/src/inbox.rs new file mode 100644 index 0000000..928173c --- /dev/null +++ b/crates/squabble-cli/src/inbox.rs @@ -0,0 +1,542 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +//! `squabble inbox-sweep` — clear notification threads whose pull request or +//! issue is already merged or closed, so the inbox holds only live items. +//! +//! Lists every thread still in the inbox (`GET /notifications?all=true`, +//! paginated), resolves each subject's state in GraphQL batches, and decides +//! per thread with `squabble_core::inbox::decide`. Only merged or closed +//! subjects are cleared; anything unread, unknown or of another kind stays. +//! +//! Dry run by default: prints what it *would* clear. `--apply` writes, per +//! thread, `DELETE /notifications/threads/{id}/subscription` (unsubscribe) +//! and then `DELETE /notifications/threads/{id}` (mark done). A thread whose +//! unsubscribe fails is not marked done. Writes stop while the REST budget +//! is at or below `--reserve`; every thread that was due to be cleared but +//! was not is listed by id — the tail is a set, never just a count. +//! +//! This needs the owner's user token with the `notifications` scope; an +//! App token cannot read notifications. +//! +//! Exit `0` = every due thread cleared (or, dry run, listed); `6` = some due +//! threads were left (budget, `--limit`, or a failed write); `2` = usage or +//! tool failure. + +use crate::board::INCOMPLETE_EXIT; +use serde_json::Value; +use squabble_core::chains::RepoId; +use squabble_core::inbox::{decide, subject_of, tally, SubjectRef, Thread, Verdict}; +use squabble_forge::inbox::{fetch_states, STATE_BATCH}; +use squabble_forge::{GhTransport, GraphQlTransport}; +use std::collections::{BTreeMap, BTreeSet}; +use std::process::{Command, ExitCode}; + +pub(crate) const USAGE: &str = "usage: squabble inbox-sweep [--owners a,b] [--repo owner/name] [--apply] [--limit N] [--reserve N]"; + +/// Owners swept when `--owners` is not given. +const DEFAULT_OWNERS: &[&str] = &["hyperpolymath", "metadatastician"]; + +/// REST requests kept back for everything else on this token. +const DEFAULT_RESERVE: u64 = 1000; + +struct Args { + owners: Vec, + repo: Option, + apply: bool, + limit: Option, + reserve: u64, +} + +/// Parse the subcommand's flags. +fn parse_args(rest: &[String]) -> Result { + let mut a = Args { + owners: DEFAULT_OWNERS.iter().map(|s| s.to_string()).collect(), + repo: None, + apply: false, + limit: None, + reserve: DEFAULT_RESERVE, + }; + let mut it = rest.iter(); + while let Some(flag) = it.next() { + let mut value = || { + it.next() + .ok_or_else(|| format!("{flag} needs a value\n{USAGE}")) + }; + match flag.as_str() { + "--owners" => { + a.owners = value()? + .split(',') + .filter(|s| !s.is_empty()) + .map(str::to_string) + .collect() + } + "--repo" => { + let r = value()?; + if r.split('/').count() != 2 || r.starts_with('/') || r.ends_with('/') { + return Err(format!("--repo wants owner/name, got `{r}`")); + } + a.repo = Some(r.clone()); + } + "--apply" => a.apply = true, + "--limit" => { + a.limit = Some( + value()? + .parse() + .map_err(|_| "--limit wants a number".to_string())?, + ) + } + "--reserve" => { + a.reserve = value()? + .parse() + .map_err(|_| "--reserve wants a number".to_string())? + } + other => return Err(format!("unknown flag `{other}`\n{USAGE}")), + } + } + if a.owners.is_empty() { + return Err("--owners is empty".into()); + } + Ok(a) +} + +/// The notification endpoints the sweep uses. A write returns the REST +/// budget left after it, when the response said. +pub(crate) trait NotificationApi { + /// Every thread still in the inbox, read or unread. + fn list(&self) -> Result, String>; + /// Unsubscribe from a thread. + fn unsubscribe(&self, id: &str) -> Result, String>; + /// Mark a thread done (removes it from the inbox). + fn mark_done(&self, id: &str) -> Result, String>; +} + +/// `gh api` with the owner's token. +struct GhNotifications; + +/// Run `gh api -i ` and return the status line, headers and body. +fn gh_with_headers(args: &[&str]) -> Result<(u16, BTreeMap, String), String> { + let out = Command::new("gh") + .arg("api") + .arg("-i") + .args(args) + .output() + .map_err(|e| format!("failed to run `gh api`: {e}"))?; + let text = String::from_utf8_lossy(&out.stdout).replace("\r\n", "\n"); + let (head, body) = text.split_once("\n\n").unwrap_or((text.as_str(), "")); + let mut lines = head.lines(); + let status = lines + .next() + .and_then(|l| l.split_whitespace().nth(1)) + .and_then(|c| c.parse().ok()) + .ok_or_else(|| { + format!( + "`gh api {}` gave no HTTP status: {}", + args.join(" "), + String::from_utf8_lossy(&out.stderr).trim() + ) + })?; + let headers = lines + .filter_map(|l| l.split_once(':')) + .map(|(k, v)| (k.trim().to_ascii_lowercase(), v.trim().to_string())) + .collect(); + Ok((status, headers, body.to_string())) +} + +/// Issue one DELETE and insist on `204 No Content`. +fn delete_expect_204(path: &str) -> Result, String> { + let (status, headers, body) = gh_with_headers(&["-X", "DELETE", path])?; + if status != 204 { + return Err(format!("DELETE {path} → HTTP {status}: {}", body.trim())); + } + Ok(headers + .get("x-ratelimit-remaining") + .and_then(|v| v.parse().ok())) +} + +/// Turn one REST notification object into a [`Thread`]. +fn thread_from_json(v: &Value) -> Option { + let s = |p: &str| v.pointer(p).and_then(Value::as_str).map(str::to_string); + Some(Thread { + id: s("/id")?, + reason: s("/reason").unwrap_or_default(), + subject_type: s("/subject/type").unwrap_or_default(), + repo: RepoId::new(s("/repository/full_name")?), + subject_url: s("/subject/url"), + unread: v.get("unread").and_then(Value::as_bool).unwrap_or(false), + updated_at: s("/updated_at").unwrap_or_default(), + }) +} + +impl NotificationApi for GhNotifications { + fn list(&self) -> Result, String> { + let out = Command::new("gh") + .args([ + "api", + "--paginate", + "--slurp", + "notifications?all=true&per_page=50", + ]) + .output() + .map_err(|e| format!("failed to run `gh api`: {e}"))?; + if !out.status.success() { + return Err(format!( + "listing notifications failed ({}): {}", + out.status, + String::from_utf8_lossy(&out.stderr).trim() + )); + } + let pages: Value = serde_json::from_slice(&out.stdout) + .map_err(|e| format!("notification listing was not JSON: {e}"))?; + let mut threads = Vec::new(); + for page in pages + .as_array() + .ok_or("notification listing was not an array of pages")? + { + for n in page + .as_array() + .ok_or("a notification page was not an array")? + { + threads.push( + thread_from_json(n).ok_or_else(|| format!("malformed notification: {n}"))?, + ); + } + } + Ok(threads) + } + + fn unsubscribe(&self, id: &str) -> Result, String> { + delete_expect_204(&format!("notifications/threads/{id}/subscription")) + } + + fn mark_done(&self, id: &str) -> Result, String> { + delete_expect_204(&format!("notifications/threads/{id}")) + } +} + +/// What a sweep found and did. +#[derive(Debug, Default)] +pub(crate) struct Sweep { + pub listed: usize, + pub out_of_scope: usize, + pub verdicts: BTreeMap, + /// Threads due to be cleared. + pub due: BTreeSet, + pub cleared: BTreeSet, + /// Due threads not cleared, with why. + pub left: BTreeMap, + pub graphql_queries: usize, +} + +/// List, decide and (with `apply`) clear. Pure apart from the two APIs. +pub(crate) fn sweep( + api: &dyn NotificationApi, + graphql: &dyn GraphQlTransport, + owners: &[String], + repo: Option<&str>, + apply: bool, + limit: Option, + reserve: u64, +) -> Result { + let all = api.list()?; + let mut out = Sweep { + listed: all.len(), + ..Sweep::default() + }; + let in_scope: Vec = all + .into_iter() + .filter(|t| { + let owner = t.repo.as_str().split('/').next().unwrap_or(""); + owners.iter().any(|o| o == owner) && repo.is_none_or(|r| t.repo.as_str() == r) + }) + .collect(); + out.out_of_scope = out.listed - in_scope.len(); + + let subjects: Vec = in_scope.iter().filter_map(subject_of).collect(); + let (states, queries) = fetch_states(graphql, &subjects, STATE_BATCH); + out.graphql_queries = queries; + + let decided: Vec<(&Thread, Verdict)> = in_scope + .iter() + .map(|t| (t, decide(t, subject_of(t).and_then(|s| states.get(&s))))) + .collect(); + out.verdicts = tally(&decided.iter().map(|(_, v)| *v).collect::>()); + let due: Vec<&Thread> = decided + .iter() + .filter(|(_, v)| v.clears()) + .map(|(t, _)| *t) + .collect(); + out.due = due.iter().map(|t| t.id.clone()).collect(); + + if !apply { + return Ok(out); + } + let mut budget: Option = None; + for t in due { + if limit.is_some_and(|l| out.cleared.len() >= l) { + out.left.insert(t.id.clone(), "--limit reached".into()); + continue; + } + if let Some(b) = budget.filter(|b| *b <= reserve) { + out.left + .insert(t.id.clone(), format!("REST budget at reserve ({b} left)")); + continue; + } + match api.unsubscribe(&t.id) { + Ok(b) => budget = b.or(budget), + Err(e) => { + out.left + .insert(t.id.clone(), format!("unsubscribe failed: {e}")); + continue; + } + } + match api.mark_done(&t.id) { + Ok(b) => { + budget = b.or(budget); + out.cleared.insert(t.id.clone()); + } + Err(e) => { + out.left + .insert(t.id.clone(), format!("unsubscribed, mark-done failed: {e}")); + } + } + } + Ok(out) +} + +/// Entry point for `squabble inbox-sweep `. +pub(crate) fn run(rest: &[String]) -> ExitCode { + let args = match parse_args(rest) { + Ok(a) => a, + Err(e) => { + eprintln!("squabble inbox-sweep: {e}"); + return ExitCode::from(2); + } + }; + let s = match sweep( + &GhNotifications, + &GhTransport, + &args.owners, + args.repo.as_deref(), + args.apply, + args.limit, + args.reserve, + ) { + Ok(s) => s, + Err(e) => { + eprintln!("squabble inbox-sweep: {e}"); + return ExitCode::from(2); + } + }; + println!( + "inbox: {} threads listed, {} out of scope, {} state queries", + s.listed, s.out_of_scope, s.graphql_queries + ); + for (v, n) in &s.verdicts { + println!(" {n:>5} {}", v.label()); + } + if !args.apply { + println!( + "dry run: {} thread(s) would be cleared; pass --apply to write", + s.due.len() + ); + return ExitCode::SUCCESS; + } + println!("cleared {} of {} due", s.cleared.len(), s.due.len()); + // The tail is the set difference, computed rather than assumed. + let unaccounted: Vec<&String> = s + .due + .iter() + .filter(|id| !s.cleared.contains(*id) && !s.left.contains_key(*id)) + .collect(); + for (id, why) in &s.left { + println!(" left {id}: {why}"); + } + for id in &unaccounted { + println!(" left {id}: not attempted (unaccounted)"); + } + if s.left.is_empty() && unaccounted.is_empty() { + ExitCode::SUCCESS + } else { + ExitCode::from(INCOMPLETE_EXIT) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + use std::cell::RefCell; + + fn t(id: &str, repo: &str, kind: &str, n: u64) -> Thread { + Thread { + id: id.into(), + reason: "author".into(), + subject_type: kind.into(), + repo: RepoId::new(repo), + subject_url: Some(format!("https://api.github.com/repos/{repo}/pulls/{n}")), + unread: false, + updated_at: "2026-10-01T00:00:00Z".into(), + } + } + + struct FakeApi { + threads: Vec, + writes: RefCell>, + fail_unsub: Option, + budget: RefCell, + } + impl NotificationApi for FakeApi { + fn list(&self) -> Result, String> { + Ok(self.threads.clone()) + } + fn unsubscribe(&self, id: &str) -> Result, String> { + if self.fail_unsub.as_deref() == Some(id) { + return Err("HTTP 500".into()); + } + self.writes.borrow_mut().push(format!("unsub {id}")); + *self.budget.borrow_mut() -= 1; + Ok(Some(*self.budget.borrow())) + } + fn mark_done(&self, id: &str) -> Result, String> { + self.writes.borrow_mut().push(format!("done {id}")); + *self.budget.borrow_mut() -= 1; + Ok(Some(*self.budget.borrow())) + } + } + + /// Answers every alias from a number → state table. + struct States(BTreeMap); + impl GraphQlTransport for States { + fn execute(&self, body: &Value) -> Result { + let vars = body["variables"].as_object().unwrap(); + let mut data = serde_json::Map::new(); + for i in 0.. { + let Some(k) = vars.get(&format!("k{i}")).and_then(Value::as_u64) else { + break; + }; + let node = self + .0 + .get(&k) + .map(|s| json!({ "__typename": "PullRequest", "prState": s })); + data.insert(format!("s{i}"), json!({ "issueOrPullRequest": node })); + } + data.insert("rateLimit".into(), json!({ "cost": 1, "remaining": 4000 })); + Ok(json!({ "data": data })) + } + } + + fn fixture(fail_unsub: Option<&str>, budget: u64) -> (FakeApi, States) { + let api = FakeApi { + threads: vec![ + t("1", "hyperpolymath/a", "PullRequest", 1), + t("2", "hyperpolymath/a", "PullRequest", 2), + t("3", "hyperpolymath/b", "PullRequest", 3), + t("4", "hyperpolymath/b", "PullRequest", 4), + t("5", "stranger/c", "PullRequest", 5), + t("6", "hyperpolymath/b", "Release", 6), + ], + writes: RefCell::new(vec![]), + fail_unsub: fail_unsub.map(str::to_string), + budget: RefCell::new(budget), + }; + // 1 merged, 2 open, 3 closed, 4 unreadable, 5 merged but out of scope. + let states = States(BTreeMap::from([ + (1, "MERGED"), + (2, "OPEN"), + (3, "CLOSED"), + (5, "MERGED"), + ])); + (api, states) + } + + fn owners() -> Vec { + vec!["hyperpolymath".into()] + } + + #[test] + fn dry_run_writes_nothing_and_lists_only_resolved_in_scope_threads() { + let (api, gql) = fixture(None, 5000); + let s = sweep(&api, &gql, &owners(), None, false, None, 1000).unwrap(); + assert!(api.writes.borrow().is_empty()); + assert_eq!(s.listed, 6); + assert_eq!(s.out_of_scope, 1); + assert_eq!(s.due, BTreeSet::from(["1".to_string(), "3".to_string()])); + assert_eq!(s.verdicts[&Verdict::KeepOpen], 1); + assert_eq!(s.verdicts[&Verdict::KeepUnknownState], 1); + assert_eq!(s.verdicts[&Verdict::KeepUnsupportedSubject], 1); + } + + #[test] + fn apply_unsubscribes_before_marking_done_and_never_touches_kept_threads() { + let (api, gql) = fixture(None, 5000); + let s = sweep(&api, &gql, &owners(), None, true, None, 1000).unwrap(); + assert_eq!( + *api.writes.borrow(), + ["unsub 1", "done 1", "unsub 3", "done 3"] + ); + assert_eq!(s.cleared, s.due); + assert!(s.left.is_empty()); + } + + #[test] + fn a_failed_unsubscribe_leaves_the_thread_in_the_inbox() { + let (api, gql) = fixture(Some("1"), 5000); + let s = sweep(&api, &gql, &owners(), None, true, None, 1000).unwrap(); + assert!(!api.writes.borrow().iter().any(|w| w == "done 1")); + assert!(s.left["1"].contains("unsubscribe failed")); + assert!(s.cleared.contains("3")); + } + + #[test] + fn writes_stop_at_the_reserve_and_the_tail_is_named() { + // Budget 1002, reserve 1000: thread 1 costs two writes, then 1000 ≤ 1000. + let (api, gql) = fixture(None, 1002); + let s = sweep(&api, &gql, &owners(), None, true, None, 1000).unwrap(); + assert_eq!(s.cleared, BTreeSet::from(["1".to_string()])); + assert!(s.left["3"].contains("reserve")); + } + + #[test] + fn limit_and_repo_narrow_the_sweep() { + let (api, gql) = fixture(None, 5000); + let s = sweep(&api, &gql, &owners(), None, true, Some(1), 1000).unwrap(); + assert_eq!(s.cleared.len(), 1); + assert!(s.left.values().all(|w| w.contains("--limit"))); + + let (api, gql) = fixture(None, 5000); + let s = sweep( + &api, + &gql, + &owners(), + Some("hyperpolymath/b"), + false, + None, + 1000, + ) + .unwrap(); + assert_eq!(s.due, BTreeSet::from(["3".to_string()])); + } + + #[test] + fn rest_notification_json_parses() { + let n = json!({ + "id": "26025256674", "reason": "author", "unread": false, + "updated_at": "2026-10-01T13:55:02Z", + "subject": { "type": "PullRequest", "url": "https://api.github.com/repos/o/r/pulls/4" }, + "repository": { "full_name": "o/r" } + }); + let th = thread_from_json(&n).unwrap(); + assert_eq!((th.id.as_str(), th.repo.as_str()), ("26025256674", "o/r")); + assert!(thread_from_json(&json!({ "id": "1" })).is_none()); + } + + #[test] + fn flags_parse_and_bad_values_fail() { + let v = |a: &[&str]| a.iter().map(|s| s.to_string()).collect::>(); + let a = parse_args(&v(&["--apply", "--limit", "3", "--repo", "o/r"])).unwrap(); + assert!(a.apply && a.limit == Some(3) && a.repo.as_deref() == Some("o/r")); + assert!(parse_args(&v(&["--repo", "nope"])).is_err()); + assert!(parse_args(&v(&["--limit"])).is_err()); + assert!(parse_args(&v(&["--frobnicate"])).is_err()); + } +} diff --git a/crates/squabble-cli/src/main.rs b/crates/squabble-cli/src/main.rs index 7050db8..0f4cc06 100644 --- a/crates/squabble-cli/src/main.rs +++ b/crates/squabble-cli/src/main.rs @@ -18,6 +18,9 @@ //! - `2` — a genuine failure: bad usage, `gh` failed, unparseable input. //! - `4` — **blocking chain finding** (`squabble chains`): a dependency cycle //! or a dead upstream. Like `3`, a reportable finding, not a tool failure. +//! - `6` — **incomplete board** (`squabble board`): rendered and published, but +//! a repo was unreadable or had more open PRs than one page; the gap is +//! stated at the top of the board. //! - `3` — **no gate**: the PR's base branch carries no `required_status_checks` //! ruleset rule, so there is nothing to triage. //! @@ -29,11 +32,13 @@ //! Note `3` says nothing about *classic* branch protection, which lives behind //! a different endpoint and is invisible to the rules API this binary queries. +mod board; #[cfg(feature = "boj")] mod boj; mod chains; mod fetch; mod fight; +mod inbox; use squabble_core::{diagnose, gate::Gate}; use std::process::ExitCode; @@ -57,6 +62,8 @@ fn main() -> ExitCode { }, Some("fight") => fight::run(&args[2..]), Some("chains") => chains::run(&args[2..]), + Some("board") => board::run(&args[2..]), + Some("inbox-sweep") => inbox::run(&args[2..]), Some("--version") | Some("-V") => { println!("squabble {}", env!("CARGO_PKG_VERSION")); ExitCode::SUCCESS @@ -71,12 +78,16 @@ fn main() -> ExitCode { squabble diagnose \n \ squabble {}\n \ squabble {}\n \ + squabble {}\n \ + squabble {}\n \ squabble --version\n\n\ EXIT CODES:\n \ - 0 ok · 2 failure · 3 no `required_status_checks` rule on the base branch · 4 chains: blocking finding\n", + 0 ok · 2 failure · 3 no `required_status_checks` rule on the base branch · 4 chains: blocking finding · 6 board or inbox-sweep: incomplete\n", env!("CARGO_PKG_VERSION"), fight::USAGE.trim_start_matches("usage: squabble "), - chains::USAGE.trim_start_matches("usage: squabble ") + chains::USAGE.trim_start_matches("usage: squabble "), + board::USAGE.trim_start_matches("usage: squabble "), + inbox::USAGE.trim_start_matches("usage: squabble ") ); ExitCode::from(2) } diff --git a/crates/squabble-core/src/board.rs b/crates/squabble-core/src/board.rs new file mode 100644 index 0000000..a1ecd9b --- /dev/null +++ b/crates/squabble-core/src/board.rs @@ -0,0 +1,626 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +//! `board` — the estate "needs me" board, pure. +//! +//! Every open pull request across the owners lands in exactly one [`Bucket`]: +//! +//! * [`Bucket::NeedsYou`] — only a human can move it (approve, merge a `CLEAN` +//! PR GitHub refuses to arm, add a gate to a repo that has none). +//! * [`Bucket::AgentWork`] — an agent can move it (conflict, red check, review +//! to answer, branch behind, automerge not yet armed). +//! * [`Bucket::Landing`] — automerge is armed and nothing is in its way; it +//! lands by itself and the board shows it only as a count. +//! * [`Bucket::Stale`] — untouched for longer than the stale threshold *and* +//! red or conflicting: the owner decides whether to close it. +//! +//! The classifier reads facts; it never fetches and never decides to merge. +//! "No required gate" is a [`Bucket::NeedsYou`] reason, never a licence to +//! arm: arming on a repo with no required contexts merges at once (AGENTS.md +//! §5c), so such a PR waits for the repo's ruleset layer instead. + +use crate::chains::RepoId; +use serde::{Deserialize, Serialize}; +use std::collections::BTreeMap; + +/// Merge-relevant facts about one repository's default branch. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RepoGate { + pub repo: RepoId, + /// Required status-check contexts in force on the default branch, summed + /// over every ruleset (and classic protection) that applies to it. + pub required_contexts: usize, + /// Highest `required_approving_review_count` among the rules that apply. + pub required_approvals: u32, +} + +/// GitHub's `MergeableState`. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +pub enum Mergeable { + Mergeable, + Conflicting, + Unknown, +} + +/// Merge-relevant facts about one open pull request. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PrFacts { + pub repo: RepoId, + pub number: u64, + pub title: String, + pub url: String, + pub author: String, + pub head_ref: String, + pub is_draft: bool, + pub mergeable: Mergeable, + /// GitHub's `MergeStateStatus` spelled as the API spells it (`CLEAN`, + /// `BLOCKED`, `BEHIND`, `DIRTY`, `UNSTABLE`, `HAS_HOOKS`, `DRAFT`, + /// `UNKNOWN`). + pub merge_state: String, + /// `APPROVED`, `CHANGES_REQUESTED`, `REVIEW_REQUIRED`, or `None`. + pub review_decision: Option, + pub auto_merge_armed: bool, + /// Head commit `statusCheckRollup.state` (`SUCCESS`, `FAILURE`, `ERROR`, + /// `PENDING`, `EXPECTED`), or `None` when the head has no checks at all. + pub rollup: Option, + /// `updatedAt`, RFC 3339. + pub updated_at: String, +} + +/// Which list a PR belongs on. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub enum Bucket { + NeedsYou, + AgentWork, + Landing, + Stale, +} + +impl Bucket { + /// Section heading on the rendered board. + pub fn heading(self) -> &'static str { + match self { + Bucket::NeedsYou => "Needs you", + Bucket::AgentWork => "Agent work", + Bucket::Landing => "Landing by itself", + Bucket::Stale => "Stale — close?", + } + } +} + +/// Why a PR is in its bucket. One reason per PR: the first that applies, in +/// the order [`classify`] tests them. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub enum Reason { + // NeedsYou + NoRequiredGate, + ApprovalRequired, + CleanMergeByHand, + // AgentWork + Draft, + Conflict, + ChangesRequested, + ChecksRed, + Behind, + ArmAutomerge, + MergeabilityUnknown, + // Landing + ArmedWaiting, + // Stale + StaleRedOrConflicting, +} + +impl Reason { + /// The bucket a reason belongs to. + pub fn bucket(self) -> Bucket { + use Reason::*; + match self { + NoRequiredGate | ApprovalRequired | CleanMergeByHand => Bucket::NeedsYou, + Draft | Conflict | ChangesRequested | ChecksRed | Behind | ArmAutomerge + | MergeabilityUnknown => Bucket::AgentWork, + ArmedWaiting => Bucket::Landing, + StaleRedOrConflicting => Bucket::Stale, + } + } + + /// One-line explanation shown on the board. + pub fn label(self) -> &'static str { + use Reason::*; + match self { + NoRequiredGate => "no required check on the default branch — add a ruleset layer (arming here would merge at once)", + ApprovalRequired => "a required approving review is missing", + CleanMergeByHand => "CLEAN: every requirement already passed, GitHub refuses to arm — merge it", + Draft => "draft", + Conflict => "merge conflict", + ChangesRequested => "changes requested by a reviewer", + ChecksRed => "a check on the head is red", + Behind => "branch is behind its base", + ArmAutomerge => "ready for automerge to be armed", + MergeabilityUnknown => "GitHub has not computed mergeability yet", + ArmedWaiting => "automerge armed, waiting on checks", + StaleRedOrConflicting => "untouched past the stale threshold and red or conflicting", + } + } +} + +/// Days since 1970-01-01 for the `YYYY-MM-DD` prefix of an RFC 3339 stamp, +/// or `None` when the prefix is not a date. +pub fn epoch_day(stamp: &str) -> Option { + let y: i64 = stamp.get(0..4)?.parse().ok()?; + let m: i64 = stamp.get(5..7)?.parse().ok()?; + let d: i64 = stamp.get(8..10)?.parse().ok()?; + if stamp.get(4..5) != Some("-") || stamp.get(7..8) != Some("-") { + return None; + } + if !(1..=12).contains(&m) || !(1..=31).contains(&d) { + return None; + } + // Howard Hinnant's days_from_civil. + let y = if m <= 2 { y - 1 } else { y }; + let era = if y >= 0 { y } else { y - 399 } / 400; + let yoe = y - era * 400; + let mp = (m + 9) % 12; + let doy = (153 * mp + 2) / 5 + d - 1; + let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy; + Some(era * 146_097 + doe - 719_468) +} + +/// RFC 3339 UTC stamp (`YYYY-MM-DDTHH:MM:SSZ`) for seconds since the epoch — +/// the inverse of [`epoch_day`] at day granularity, so no date crate is needed. +pub fn rfc3339_from_unix(secs: u64) -> String { + let days = (secs / 86_400) as i64; + let rem = secs % 86_400; + // Howard Hinnant's civil_from_days. + let z = days + 719_468; + let era = z.div_euclid(146_097); + let doe = z - era * 146_097; + let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; + let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); + let mp = (5 * doy + 2) / 153; + let d = doy - (153 * mp + 2) / 5 + 1; + let m = if mp < 10 { mp + 3 } else { mp - 9 }; + let y = yoe + era * 400 + i64::from(m <= 2); + format!( + "{y:04}-{m:02}-{d:02}T{:02}:{:02}:{:02}Z", + rem / 3600, + rem % 3600 / 60, + rem % 60 + ) +} + +/// Place one PR. `today` is [`epoch_day`] of the run; `stale_days` is the +/// threshold past which a red or conflicting PR is offered for closing. +/// +/// A missing gate outranks everything except staleness: whatever else is +/// wrong with such a PR, the repo has to grow a gate before automation can +/// land anything there. +pub fn classify(pr: &PrFacts, gate: &RepoGate, today: i64, stale_days: i64) -> Reason { + let red = matches!(pr.rollup.as_deref(), Some("FAILURE" | "ERROR")); + let conflicting = pr.mergeable == Mergeable::Conflicting; + let age = epoch_day(&pr.updated_at).map(|d| today - d); + if (red || conflicting) && age.is_some_and(|a| a > stale_days) { + return Reason::StaleRedOrConflicting; + } + if gate.required_contexts == 0 { + return Reason::NoRequiredGate; + } + if pr.is_draft { + return Reason::Draft; + } + if conflicting { + return Reason::Conflict; + } + match pr.review_decision.as_deref() { + Some("CHANGES_REQUESTED") => return Reason::ChangesRequested, + Some("REVIEW_REQUIRED") if gate.required_approvals > 0 => return Reason::ApprovalRequired, + _ => {} + } + if red { + return Reason::ChecksRed; + } + if pr.merge_state == "BEHIND" { + return Reason::Behind; + } + if pr.auto_merge_armed { + return Reason::ArmedWaiting; + } + match (pr.mergeable, pr.merge_state.as_str()) { + (_, "CLEAN") => Reason::CleanMergeByHand, + (Mergeable::Unknown, _) | (_, "UNKNOWN") => Reason::MergeabilityUnknown, + _ => Reason::ArmAutomerge, + } +} + +/// One run of the board: what was read, what could not be, and every PR placed. +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct Board { + /// RFC 3339 time of the run. + pub generated_at: String, + pub owners: Vec, + /// Repositories enumerated (non-archived), per owner. + pub repos_enumerated: BTreeMap, + /// Open PRs the enumeration reported, per owner — the denominator. + pub open_prs_reported: BTreeMap, + /// Repositories that could not be read, with the reason. Never dropped. + pub unavailable: Vec<(RepoId, String)>, + /// Repositories whose open PRs exceeded one page; the rest were not read. + pub truncated: Vec, + pub placed: Vec<(PrFacts, Reason)>, +} + +impl Board { + /// PRs actually placed, per owner — compared against + /// [`Board::open_prs_reported`] so a shortfall is visible, never silent. + pub fn placed_per_owner(&self) -> BTreeMap { + let mut m = BTreeMap::new(); + for (pr, _) in &self.placed { + let owner = pr.repo.as_str().split('/').next().unwrap_or("").to_string(); + *m.entry(owner).or_insert(0) += 1; + } + m + } +} + +/// GitHub rejects an issue body over 65 536 characters. +pub const ISSUE_BODY_LIMIT: usize = 65_536; + +/// Render the board as GitHub-flavoured Markdown no longer than `max_bytes`. +/// +/// `Needs you` is listed in full first; `Agent work` and `Stale` are grouped +/// by reason inside collapsed sections; `Landing` is a count. When the text +/// would exceed `max_bytes`, the longest lists are cut and each cut says how +/// many rows it hid — the totals in the header are always complete. +pub fn render_markdown(board: &Board, max_bytes: usize) -> String { + let mut per_page = 400usize; + loop { + let text = render_with_cap(board, per_page); + if text.len() <= max_bytes || per_page == 0 { + return text; + } + per_page /= 2; + } +} + +/// Render the whole board with each detail list cut at `cap` lines. +fn render_with_cap(board: &Board, cap: usize) -> String { + let mut by_reason: BTreeMap> = BTreeMap::new(); + for (pr, r) in &board.placed { + by_reason.entry(*r).or_default().push(pr); + } + let count = |b: Bucket| -> usize { + by_reason + .iter() + .filter(|(r, _)| r.bucket() == b) + .map(|(_, v)| v.len()) + .sum() + }; + + let mut s = String::new(); + s.push_str("# Estate: needs me\n\n"); + s.push_str(&format!( + "_Generated {} by `squabble board` for {}. This body is rewritten on every run; it never comments, so it never notifies._\n\n", + board.generated_at, + board.owners.join(", ") + )); + s.push_str("| | count |\n|---|---:|\n"); + for b in [ + Bucket::NeedsYou, + Bucket::AgentWork, + Bucket::Landing, + Bucket::Stale, + ] { + s.push_str(&format!("| **{}** | {} |\n", b.heading(), count(b))); + } + s.push('\n'); + + let placed = board.placed_per_owner(); + s.push_str("**Coverage** — "); + let cov: Vec = board + .owners + .iter() + .map(|o| { + format!( + "{o}: {} repos, {} open PRs reported, {} placed", + board.repos_enumerated.get(o).copied().unwrap_or(0), + board.open_prs_reported.get(o).copied().unwrap_or(0), + placed.get(o).copied().unwrap_or(0) + ) + }) + .collect(); + s.push_str(&cov.join(" · ")); + s.push('\n'); + if !board.unavailable.is_empty() || !board.truncated.is_empty() { + s.push_str(&format!( + "\n> ⚠ **Incomplete:** {} repo(s) unreadable, {} repo(s) with more open PRs than one page. Their PRs are missing from the lists below.\n", + board.unavailable.len(), + board.truncated.len() + )); + for (r, why) in &board.unavailable { + s.push_str(&format!("> - `{}` — {}\n", r.as_str(), why)); + } + for r in &board.truncated { + s.push_str(&format!("> - `{}` — truncated at one page\n", r.as_str())); + } + } + + // Needs you: no-gate PRs collapse to one line per repo (the fix is per + // repo); the rest are listed one per PR. + s.push_str(&format!("\n## {}\n", Bucket::NeedsYou.heading())); + let mut any = false; + for (r, prs) in by_reason + .iter() + .filter(|(r, _)| r.bucket() == Bucket::NeedsYou) + { + any = true; + s.push_str(&format!("\n### {} ({})\n", r.label(), prs.len())); + if *r == Reason::NoRequiredGate { + let mut repos: BTreeMap<&str, usize> = BTreeMap::new(); + for p in prs { + *repos.entry(p.repo.as_str()).or_insert(0) += 1; + } + push_capped(&mut s, repos.iter(), cap, |(repo, n)| { + format!("- `{repo}` — {n} open PR(s)") + }); + } else { + push_capped(&mut s, prs.iter(), cap, |p| pr_line(p)); + } + } + if !any { + s.push_str("\nNothing. 🎉\n"); + } + + for b in [Bucket::AgentWork, Bucket::Stale] { + s.push_str(&format!("\n## {} ({})\n", b.heading(), count(b))); + for (r, prs) in by_reason.iter().filter(|(r, _)| r.bucket() == b) { + s.push_str(&format!( + "\n
{} — {}\n\n", + r.label(), + prs.len() + )); + push_capped(&mut s, prs.iter(), cap, |p| pr_line(p)); + s.push_str("\n
\n"); + } + } + + s.push_str(&format!( + "\n## {} ({})\n\nArmed and unobstructed; these need nothing from anyone.\n", + Bucket::Landing.heading(), + count(Bucket::Landing) + )); + s +} + +/// One Markdown list line for a pull request. +fn pr_line(p: &PrFacts) -> String { + let title: String = p.title.replace('|', "\\|").chars().take(90).collect(); + format!( + "- [{}#{}]({}) {} — @{}", + p.repo.as_str(), + p.number, + p.url, + title, + p.author + ) +} + +/// Append up to `cap` lines, then an "…and N more" line for the rest. +fn push_capped( + s: &mut String, + items: impl ExactSizeIterator, + cap: usize, + line: impl Fn(T) -> String, +) { + let total = items.len(); + for item in items.take(cap) { + s.push_str(&line(item)); + s.push('\n'); + } + if total > cap { + s.push_str(&format!( + "- …and {} more (run `squabble board` locally for the full list)\n", + total - cap + )); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn gate(ctx: usize, approvals: u32) -> RepoGate { + RepoGate { + repo: RepoId::new("o/r"), + required_contexts: ctx, + required_approvals: approvals, + } + } + + fn pr() -> PrFacts { + PrFacts { + repo: RepoId::new("o/r"), + number: 1, + title: "t".into(), + url: "https://github.com/o/r/pull/1".into(), + author: "a".into(), + head_ref: "h".into(), + is_draft: false, + mergeable: Mergeable::Mergeable, + merge_state: "BLOCKED".into(), + review_decision: None, + auto_merge_armed: false, + rollup: Some("PENDING".into()), + updated_at: "2026-10-01T10:00:00Z".into(), + } + } + + const TODAY: i64 = 20_727; // 2026-10-01 + + #[test] + fn epoch_day_matches_known_dates() { + assert_eq!(epoch_day("1970-01-01T00:00:00Z"), Some(0)); + assert_eq!(epoch_day("2000-03-01"), Some(11_017)); + assert_eq!(epoch_day("2026-10-01T13:00:00Z"), Some(TODAY)); + assert_eq!(epoch_day("not a date"), None); + assert_eq!(epoch_day("2026-13-01"), None); + } + + #[test] + fn rfc3339_round_trips_through_epoch_day() { + assert_eq!(rfc3339_from_unix(0), "1970-01-01T00:00:00Z"); + assert_eq!(rfc3339_from_unix(951_782_400), "2000-02-29T00:00:00Z"); + for day in [0i64, 11_016, 11_017, TODAY, 40_000] { + let s = rfc3339_from_unix(day as u64 * 86_400 + 3_723); + assert_eq!(epoch_day(&s), Some(day), "{s}"); + assert!(s.ends_with("T01:02:03Z"), "{s}"); + } + } + + #[test] + fn no_gate_is_never_armable() { + // The §5c trap: a PR that would otherwise be "arm it" must not be, + // because there is nothing for automerge to wait on. + assert_eq!( + classify(&pr(), &gate(0, 0), TODAY, 30), + Reason::NoRequiredGate + ); + let mut clean = pr(); + clean.merge_state = "CLEAN".into(); + assert_eq!( + classify(&clean, &gate(0, 0), TODAY, 30), + Reason::NoRequiredGate + ); + } + + #[test] + fn gated_unarmed_blocked_pr_is_armable() { + assert_eq!( + classify(&pr(), &gate(3, 0), TODAY, 30), + Reason::ArmAutomerge + ); + } + + #[test] + fn armed_and_unobstructed_is_landing() { + let mut p = pr(); + p.auto_merge_armed = true; + assert_eq!(classify(&p, &gate(3, 0), TODAY, 30), Reason::ArmedWaiting); + assert_eq!(Reason::ArmedWaiting.bucket(), Bucket::Landing); + } + + #[test] + fn armed_but_red_is_agent_work_not_landing() { + let mut p = pr(); + p.auto_merge_armed = true; + p.rollup = Some("FAILURE".into()); + assert_eq!(classify(&p, &gate(3, 0), TODAY, 30), Reason::ChecksRed); + } + + #[test] + fn clean_unarmed_needs_a_human_merge() { + let mut p = pr(); + p.merge_state = "CLEAN".into(); + assert_eq!( + classify(&p, &gate(3, 0), TODAY, 30), + Reason::CleanMergeByHand + ); + assert_eq!(Reason::CleanMergeByHand.bucket(), Bucket::NeedsYou); + } + + #[test] + fn approval_only_needs_you_when_the_rules_require_one() { + let mut p = pr(); + p.review_decision = Some("REVIEW_REQUIRED".into()); + assert_eq!( + classify(&p, &gate(3, 1), TODAY, 30), + Reason::ApprovalRequired + ); + assert_eq!(classify(&p, &gate(3, 0), TODAY, 30), Reason::ArmAutomerge); + } + + #[test] + fn conflicts_and_drafts_are_agent_work() { + let mut p = pr(); + p.mergeable = Mergeable::Conflicting; + assert_eq!(classify(&p, &gate(3, 0), TODAY, 30), Reason::Conflict); + let mut d = pr(); + d.is_draft = true; + assert_eq!(classify(&d, &gate(3, 0), TODAY, 30), Reason::Draft); + } + + #[test] + fn old_and_red_is_stale_but_old_and_green_is_not() { + let mut p = pr(); + p.updated_at = "2026-08-01T00:00:00Z".into(); + p.rollup = Some("FAILURE".into()); + assert_eq!( + classify(&p, &gate(0, 0), TODAY, 30), + Reason::StaleRedOrConflicting + ); + p.rollup = Some("SUCCESS".into()); + assert_eq!(classify(&p, &gate(3, 0), TODAY, 30), Reason::ArmAutomerge); + } + + #[test] + fn every_reason_maps_to_the_bucket_its_section_claims() { + use Reason::*; + for r in [NoRequiredGate, ApprovalRequired, CleanMergeByHand] { + assert_eq!(r.bucket(), Bucket::NeedsYou); + } + for r in [ + Draft, + Conflict, + ChangesRequested, + ChecksRed, + Behind, + ArmAutomerge, + MergeabilityUnknown, + ] { + assert_eq!(r.bucket(), Bucket::AgentWork); + } + } + + fn big_board(n: usize) -> Board { + let mut b = Board { + generated_at: "2026-10-01T13:00:00Z".into(), + owners: vec!["o".into()], + ..Board::default() + }; + b.repos_enumerated.insert("o".into(), 1); + b.open_prs_reported.insert("o".into(), n); + for i in 0..n { + let mut p = pr(); + p.number = i as u64; + p.title = "x".repeat(80); + b.placed.push((p, Reason::Conflict)); + } + b + } + + #[test] + fn render_respects_the_issue_body_limit_and_says_what_it_hid() { + let b = big_board(5_000); + let md = render_markdown(&b, ISSUE_BODY_LIMIT); + assert!(md.len() <= ISSUE_BODY_LIMIT, "{} bytes", md.len()); + assert!(md.contains("more (run `squabble board` locally")); + // The header totals are never cut. + assert!(md.contains("| **Agent work** | 5000 |")); + assert!(md.contains("5000 open PRs reported, 5000 placed")); + } + + #[test] + fn render_lists_everything_when_it_fits() { + let b = big_board(3); + let md = render_markdown(&b, ISSUE_BODY_LIMIT); + assert!(!md.contains("more (run")); + assert_eq!(md.matches("](https://github.com/o/r/pull/1)").count(), 3); + } + + #[test] + fn unreadable_repos_are_announced_not_dropped() { + let mut b = big_board(1); + b.unavailable + .push((RepoId::new("o/gone"), "rate limited".into())); + let md = render_markdown(&b, ISSUE_BODY_LIMIT); + assert!(md.contains("Incomplete")); + assert!(md.contains("`o/gone` — rate limited")); + } +} diff --git a/crates/squabble-core/src/inbox.rs b/crates/squabble-core/src/inbox.rs new file mode 100644 index 0000000..190087b --- /dev/null +++ b/crates/squabble-core/src/inbox.rs @@ -0,0 +1,225 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +//! `inbox` — which notification threads are finished, pure. +//! +//! A thread whose subject (a pull request or an issue) is merged or closed +//! asks nothing more of anyone: `squabble inbox-sweep` unsubscribes from it +//! and marks it done. Every other thread is kept, with the reason, so the +//! inbox ends up holding only what is still live. +//! +//! Nothing here is guessed. A subject whose state could not be read is kept; +//! a subject type the sweep does not understand (a release, a check suite, a +//! discussion) is kept. Unsubscribing is reversible — GitHub re-subscribes +//! on a new mention or review request — but a wrongly cleared thread is an +//! item the owner never sees, so the default is always "keep". + +use crate::chains::RepoId; +use serde::{Deserialize, Serialize}; +use std::collections::BTreeMap; + +/// One notification thread, as the REST `notifications` endpoint lists it. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct Thread { + pub id: String, + pub reason: String, + /// `subject.type`: `PullRequest`, `Issue`, `Release`, `CheckSuite`, … + pub subject_type: String, + pub repo: RepoId, + /// `subject.url`, an API URL; `None` for subjects without one. + pub subject_url: Option, + pub unread: bool, + pub updated_at: String, +} + +/// The issue or pull request a thread is about. +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub struct SubjectRef { + pub repo: RepoId, + pub number: u64, +} + +impl SubjectRef { + /// Parse `https://api.github.com/repos/{o}/{r}/(pulls|issues)/{n}`. + /// Anything else — a different host, a commit URL, a release — is `None`. + pub fn from_api_url(url: &str) -> Option { + let rest = url.strip_prefix("https://api.github.com/repos/")?; + let parts: Vec<&str> = rest.split('/').collect(); + let [owner, name, kind, number] = parts.as_slice() else { + return None; + }; + if owner.is_empty() || name.is_empty() || !matches!(*kind, "pulls" | "issues") { + return None; + } + Some(Self { + repo: RepoId::new(format!("{owner}/{name}")), + number: number.parse().ok()?, + }) + } +} + +/// A subject's state as the forge reported it. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub enum SubjectState { + Open, + Merged, + Closed, + /// Could not be read; the reason is carried, and the thread is kept. + Unknown(String), +} + +/// What to do with one thread. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub enum Verdict { + /// Unsubscribe, then mark done. + ClearMerged, + ClearClosed, + KeepOpen, + KeepUnknownState, + KeepUnsupportedSubject, +} + +impl Verdict { + /// True for the verdicts that write. + pub fn clears(self) -> bool { + matches!(self, Verdict::ClearMerged | Verdict::ClearClosed) + } + + /// One-line label for reports. + pub fn label(self) -> &'static str { + match self { + Verdict::ClearMerged => "clear: subject merged", + Verdict::ClearClosed => "clear: subject closed", + Verdict::KeepOpen => "keep: subject still open", + Verdict::KeepUnknownState => "keep: subject state could not be read", + Verdict::KeepUnsupportedSubject => "keep: not a pull request or issue", + } + } +} + +/// The subject a thread refers to, when the sweep understands it. +pub fn subject_of(t: &Thread) -> Option { + if !matches!(t.subject_type.as_str(), "PullRequest" | "Issue") { + return None; + } + SubjectRef::from_api_url(t.subject_url.as_deref()?) +} + +/// Decide one thread given the state lookup (`None` = not looked up). +pub fn decide(t: &Thread, state: Option<&SubjectState>) -> Verdict { + if subject_of(t).is_none() { + return Verdict::KeepUnsupportedSubject; + } + match state { + Some(SubjectState::Merged) => Verdict::ClearMerged, + Some(SubjectState::Closed) => Verdict::ClearClosed, + Some(SubjectState::Open) => Verdict::KeepOpen, + Some(SubjectState::Unknown(_)) | None => Verdict::KeepUnknownState, + } +} + +/// Count verdicts, for the before-anything-is-written summary. +pub fn tally(verdicts: &[Verdict]) -> BTreeMap { + let mut m = BTreeMap::new(); + for v in verdicts { + *m.entry(*v).or_insert(0) += 1; + } + m +} + +#[cfg(test)] +mod tests { + use super::*; + + fn thread(kind: &str, url: Option<&str>) -> Thread { + Thread { + id: "1".into(), + reason: "author".into(), + subject_type: kind.into(), + repo: RepoId::new("o/r"), + subject_url: url.map(str::to_string), + unread: false, + updated_at: "2026-10-01T00:00:00Z".into(), + } + } + + #[test] + fn subject_urls_parse_and_reject() { + let s = SubjectRef::from_api_url("https://api.github.com/repos/o/r/pulls/12").unwrap(); + assert_eq!((s.repo.as_str(), s.number), ("o/r", 12)); + assert!(SubjectRef::from_api_url("https://api.github.com/repos/o/r/issues/3").is_some()); + for bad in [ + "https://api.github.com/repos/o/r/commits/abc", + "https://api.github.com/repos/o/r/pulls/x", + "https://api.github.com/repos/o/r/pulls/1/extra", + "https://evil.example/repos/o/r/pulls/1", + "https://api.github.com/repos//r/pulls/1", + ] { + assert_eq!(SubjectRef::from_api_url(bad), None, "{bad}"); + } + } + + #[test] + fn only_resolved_subjects_clear() { + let t = thread( + "PullRequest", + Some("https://api.github.com/repos/o/r/pulls/1"), + ); + assert_eq!( + decide(&t, Some(&SubjectState::Merged)), + Verdict::ClearMerged + ); + assert_eq!( + decide(&t, Some(&SubjectState::Closed)), + Verdict::ClearClosed + ); + assert_eq!(decide(&t, Some(&SubjectState::Open)), Verdict::KeepOpen); + assert!(!Verdict::KeepOpen.clears()); + } + + #[test] + fn unread_state_and_unknown_kinds_are_kept() { + let t = thread( + "PullRequest", + Some("https://api.github.com/repos/o/r/pulls/1"), + ); + assert_eq!(decide(&t, None), Verdict::KeepUnknownState); + assert_eq!( + decide(&t, Some(&SubjectState::Unknown("502".into()))), + Verdict::KeepUnknownState + ); + let rel = thread( + "Release", + Some("https://api.github.com/repos/o/r/releases/9"), + ); + assert_eq!( + decide(&rel, Some(&SubjectState::Closed)), + Verdict::KeepUnsupportedSubject + ); + let none = thread("PullRequest", None); + assert_eq!( + decide(&none, Some(&SubjectState::Merged)), + Verdict::KeepUnsupportedSubject + ); + } + + #[test] + fn a_pull_request_thread_with_an_issue_url_is_still_understood() { + // GitHub sometimes points PR threads at the issues endpoint. + let t = thread( + "PullRequest", + Some("https://api.github.com/repos/o/r/issues/5"), + ); + assert_eq!(subject_of(&t).unwrap().number, 5); + } + + #[test] + fn tally_counts_each_verdict() { + let m = tally(&[ + Verdict::ClearMerged, + Verdict::ClearMerged, + Verdict::KeepOpen, + ]); + assert_eq!(m[&Verdict::ClearMerged], 2); + assert_eq!(m[&Verdict::KeepOpen], 1); + } +} diff --git a/crates/squabble-core/src/lib.rs b/crates/squabble-core/src/lib.rs index 101d100..5508d56 100644 --- a/crates/squabble-core/src/lib.rs +++ b/crates/squabble-core/src/lib.rs @@ -15,8 +15,10 @@ //! detachability is the whole point. pub mod admission; +pub mod board; pub mod chains; pub mod gate; +pub mod inbox; pub mod moves; pub mod outcome; pub mod polarity; diff --git a/crates/squabble-forge/graphql/board_repo.graphql b/crates/squabble-forge/graphql/board_repo.graphql new file mode 100644 index 0000000..6f1a0a2 --- /dev/null +++ b/crates/squabble-forge/graphql/board_repo.graphql @@ -0,0 +1,53 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +# +# One repository's merge gate and open pull requests, for `squabble board`. +# `board_query` in src/board.rs aliases this fragment once per repo. +# +# The gate is read from BOTH sources GitHub enforces: `rules` (rulesets, the +# effective set for the default branch) and `branchProtectionRule` (classic +# protection), because `rules` does not report classic protection. +fragment BoardRepo on Repository { + nameWithOwner + defaultBranchRef { + name + rules(first: 100) { + nodes { + type + parameters { + ... on RequiredStatusChecksParameters { + requiredStatusChecks { context } + } + ... on PullRequestParameters { + requiredApprovingReviewCount + } + } + } + } + branchProtectionRule { + requiresStatusChecks + requiredStatusCheckContexts + requiresApprovingReviews + requiredApprovingReviewCount + } + } + pullRequests(states: [OPEN], first: 100, orderBy: { field: UPDATED_AT, direction: DESC }) { + totalCount + nodes { + number + title + url + isDraft + mergeable + mergeStateStatus + reviewDecision + updatedAt + headRefName + author { login } + autoMergeRequest { enabledAt } + commits(last: 1) { + nodes { commit { statusCheckRollup { state } } } + } + } + } +} diff --git a/crates/squabble-forge/graphql/estate_repos.graphql b/crates/squabble-forge/graphql/estate_repos.graphql new file mode 100644 index 0000000..f0fcd46 --- /dev/null +++ b/crates/squabble-forge/graphql/estate_repos.graphql @@ -0,0 +1,21 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +# +# One page of an owner's non-archived repositories, with each one's open-PR +# total. `squabble board` pages through this to build its denominator (repos +# enumerated, open PRs reported) before reading any PR in detail. +# ownerAffiliations is pinned to OWNER: the default also lists collaborator +# repos belonging to someone else. +query EstateRepos($login: String!, $after: String) { + rateLimit { cost remaining resetAt } + repositoryOwner(login: $login) { + repositories(first: 100, after: $after, isArchived: false, ownerAffiliations: [OWNER]) { + totalCount + pageInfo { hasNextPage endCursor } + nodes { + nameWithOwner + pullRequests(states: [OPEN]) { totalCount } + } + } + } +} diff --git a/crates/squabble-forge/src/board.rs b/crates/squabble-forge/src/board.rs new file mode 100644 index 0000000..99dc2ee --- /dev/null +++ b/crates/squabble-forge/src/board.rs @@ -0,0 +1,730 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +//! Forge reads for `squabble board`: enumerate an owner's repositories, then +//! read each repo's merge gate and open pull requests in batched queries. +//! +//! Same contract as the crate root: variables never interpolated, every +//! document schema-validated in tests, and fail-closed — a repo that could not +//! be read is reported with its reason, never silently dropped, and a short +//! enumeration (fewer repos listed than `totalCount`) is an error. + +use crate::{split_slug, GraphQlTransport, RateSample}; +use serde_json::{json, Value}; +use squabble_core::board::{Mergeable, PrFacts, RepoGate}; +use squabble_core::chains::RepoId; +use std::collections::BTreeSet; + +/// Paged repository enumeration for one owner. +pub const ESTATE_REPOS: &str = include_str!("../graphql/estate_repos.graphql"); +/// Per-repo gate + open PRs, aliased once per repo by [`board_query`]. +pub const BOARD_REPO: &str = include_str!("../graphql/board_repo.graphql"); + +/// Repos per board query. PR nodes are heavier than workflow trees, so this +/// is smaller than [`crate::DEFAULT_BATCH`]. +pub const BOARD_BATCH: usize = 10; + +/// GitHub's page size for the PR connection in [`BOARD_REPO`]. +pub const PR_PAGE: usize = 100; + +/// Every non-archived repo an owner owns, with its open-PR total. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct OwnerListing { + pub login: String, + pub repos: Vec<(RepoId, usize)>, + pub rate: Option, +} + +impl OwnerListing { + /// Sum of the per-repo open-PR totals — the board's denominator. + pub fn open_prs(&self) -> usize { + self.repos.iter().map(|(_, n)| n).sum() + } +} + +/// The `rateLimit` block of a response, when it carried one. +fn rate_of(resp: &Value) -> Option { + let r = resp.pointer("/data/rateLimit")?; + Some(RateSample { + cost: r.get("cost")?.as_u64()?, + remaining: r.get("remaining")?.as_u64()?, + reset_at: r.get("resetAt").and_then(Value::as_str).map(str::to_string), + }) +} + +/// A one-line reason for a response that carried no usable data. +fn error_text(resp: &Value) -> String { + let msgs: Vec<&str> = resp + .get("errors") + .and_then(Value::as_array) + .into_iter() + .flatten() + .filter_map(|e| e.get("message").and_then(Value::as_str)) + .collect(); + if msgs.is_empty() { + // A REST-shaped error body (`{"message": …}`) arrives from gateway + // failures; say so rather than "no data". + resp.get("message").and_then(Value::as_str).map_or_else( + || format!("response carried no data: {}", short(resp)), + str::to_string, + ) + } else { + msgs.join("; ") + } +} + +/// First 200 characters of a response, for an error message. +fn short(v: &Value) -> String { + v.to_string().chars().take(200).collect() +} + +/// Attempts per listing page. A gateway 502 on one page would otherwise +/// abort the whole board. +const LISTING_ATTEMPTS: u32 = 3; + +/// Execute `body`, retrying only a transport failure (no JSON came back) +/// with 1 s, 2 s, … back-off. A GraphQL error is an answer, not retried. +fn execute_with_retry( + transport: &dyn GraphQlTransport, + body: &Value, + attempts: u32, +) -> Result { + let mut last = String::new(); + for n in 0..attempts.max(1) { + if n > 0 { + std::thread::sleep(std::time::Duration::from_secs(1 << (n - 1))); + } + match transport.execute(body) { + Ok(v) => return Ok(v), + Err(e) => last = e, + } + } + Err(format!("{last} (after {} attempts)", attempts.max(1))) +} + +/// List every non-archived repository `login` owns, following cursors to the +/// end. Errs on any page that carries no data, and on a listing whose length +/// disagrees with GitHub's own `totalCount`. +pub fn enumerate_owner( + transport: &dyn GraphQlTransport, + login: &str, +) -> Result { + let mut out = OwnerListing { + login: login.to_string(), + ..OwnerListing::default() + }; + let mut after: Option = None; + let mut total: Option = None; + loop { + let body = json!({ + "query": ESTATE_REPOS, + "variables": { "login": login, "after": after }, + }); + let resp = execute_with_retry(transport, &body, LISTING_ATTEMPTS)?; + out.rate = rate_of(&resp).or(out.rate); + let conn = resp + .pointer("/data/repositoryOwner/repositories") + .filter(|c| c.is_object()) + .ok_or_else(|| format!("listing `{login}`: {}", error_text(&resp)))?; + total = conn.get("totalCount").and_then(Value::as_u64).or(total); + for n in conn + .get("nodes") + .and_then(Value::as_array) + .into_iter() + .flatten() + { + let Some(slug) = n.get("nameWithOwner").and_then(Value::as_str) else { + return Err(format!("listing `{login}`: a repository node had no name")); + }; + let open = n + .pointer("/pullRequests/totalCount") + .and_then(Value::as_u64) + .ok_or_else(|| format!("listing `{login}`: `{slug}` had no PR count"))?; + out.repos.push((RepoId::new(slug), open as usize)); + } + let next = conn + .pointer("/pageInfo/hasNextPage") + .and_then(Value::as_bool) + .unwrap_or(false); + after = conn + .pointer("/pageInfo/endCursor") + .and_then(Value::as_str) + .map(str::to_string); + if !next || after.is_none() { + break; + } + } + match total { + Some(t) if t as usize == out.repos.len() => Ok(out), + Some(t) => Err(format!( + "listing `{login}`: GitHub reports {t} repositories but {} were listed", + out.repos.len() + )), + None => Err(format!("listing `{login}`: no totalCount in the response")), + } +} + +/// Build one batched board query: rate probe plus [`BOARD_REPO`] aliased +/// `r0…rN`, owner/name passed as variables. +pub fn board_query(batch: &[RepoId]) -> (String, Value) { + let mut params = Vec::new(); + let mut fields = Vec::new(); + let mut vars = serde_json::Map::new(); + for (i, repo) in batch.iter().enumerate() { + let (o, n) = split_slug(repo).unwrap_or(("", "")); + params.push(format!("$o{i}: String!, $n{i}: String!")); + fields.push(format!( + " r{i}: repository(owner: $o{i}, name: $n{i}) {{ ...BoardRepo }}" + )); + vars.insert(format!("o{i}"), json!(o)); + vars.insert(format!("n{i}"), json!(n)); + } + let doc = format!( + "query EstateBoard({}) {{\n rateLimit {{ cost remaining resetAt }}\n{}\n}}\n{}", + params.join(", "), + fields.join("\n"), + BOARD_REPO + ); + (doc, Value::Object(vars)) +} + +/// What one repo's board read produced. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum RepoRead { + Read { + gate: RepoGate, + prs: Vec, + /// More open PRs exist than one page returned. + truncated: bool, + }, + Unavailable(String), +} + +/// Turn one batched board response into reads, in `batch` order. +pub fn parse_board_response(batch: &[RepoId], resp: &Value) -> (Vec, Option) { + let rate = rate_of(resp); + let data = resp.get("data").filter(|d| d.is_object()); + let alias_error = |alias: &str| -> Option { + resp.get("errors") + .and_then(Value::as_array)? + .iter() + .find(|e| e.pointer("/path/0").and_then(Value::as_str) == Some(alias)) + .map(|e| { + e.get("message") + .and_then(Value::as_str) + .unwrap_or("GraphQL error without a message") + .to_string() + }) + }; + let reads = batch + .iter() + .enumerate() + .map(|(i, repo)| { + let alias = format!("r{i}"); + let Some(data) = data else { + return RepoRead::Unavailable(error_text(resp)); + }; + match data.get(&alias) { + Some(node) if node.is_object() => parse_board_repo(repo, node), + _ => RepoRead::Unavailable( + alias_error(&alias) + .unwrap_or_else(|| "null repository with no matching error".into()), + ), + } + }) + .collect(); + (reads, rate) +} + +/// Read one repository alias: its merge gate (rulesets ∪ classic protection) and open PRs. +fn parse_board_repo(repo: &RepoId, node: &Value) -> RepoRead { + let mut contexts: BTreeSet = BTreeSet::new(); + let mut approvals: u32 = 0; + let branch = node.get("defaultBranchRef").filter(|b| b.is_object()); + for rule in branch + .and_then(|b| b.pointer("/rules/nodes")) + .and_then(Value::as_array) + .into_iter() + .flatten() + { + let params = rule.get("parameters"); + match rule.get("type").and_then(Value::as_str) { + Some("REQUIRED_STATUS_CHECKS") => { + for c in params + .and_then(|p| p.get("requiredStatusChecks")) + .and_then(Value::as_array) + .into_iter() + .flatten() + { + if let Some(ctx) = c.get("context").and_then(Value::as_str) { + contexts.insert(ctx.to_string()); + } + } + } + Some("PULL_REQUEST") => { + let n = params + .and_then(|p| p.get("requiredApprovingReviewCount")) + .and_then(Value::as_u64) + .unwrap_or(0) as u32; + approvals = approvals.max(n); + } + _ => {} + } + } + if let Some(bp) = branch + .and_then(|b| b.get("branchProtectionRule")) + .filter(|b| b.is_object()) + { + if bp.get("requiresStatusChecks").and_then(Value::as_bool) == Some(true) { + for c in bp + .get("requiredStatusCheckContexts") + .and_then(Value::as_array) + .into_iter() + .flatten() + .filter_map(Value::as_str) + { + contexts.insert(c.to_string()); + } + } + if bp.get("requiresApprovingReviews").and_then(Value::as_bool) == Some(true) { + let n = bp + .get("requiredApprovingReviewCount") + .and_then(Value::as_u64) + .unwrap_or(0) as u32; + approvals = approvals.max(n); + } + } + let gate = RepoGate { + repo: repo.clone(), + required_contexts: contexts.len(), + required_approvals: approvals, + }; + + let conn = node.get("pullRequests"); + let total = conn + .and_then(|c| c.get("totalCount")) + .and_then(Value::as_u64) + .unwrap_or(0) as usize; + let mut prs = Vec::new(); + for p in conn + .and_then(|c| c.get("nodes")) + .and_then(Value::as_array) + .into_iter() + .flatten() + { + let s = |k: &str| { + p.get(k) + .and_then(Value::as_str) + .unwrap_or_default() + .to_string() + }; + let Some(number) = p.get("number").and_then(Value::as_u64) else { + return RepoRead::Unavailable("a pull request node had no number".into()); + }; + prs.push(PrFacts { + repo: repo.clone(), + number, + title: s("title"), + url: s("url"), + author: p + .pointer("/author/login") + .and_then(Value::as_str) + .unwrap_or("ghost") + .to_string(), + head_ref: s("headRefName"), + is_draft: p.get("isDraft").and_then(Value::as_bool).unwrap_or(false), + mergeable: match p.get("mergeable").and_then(Value::as_str) { + Some("MERGEABLE") => Mergeable::Mergeable, + Some("CONFLICTING") => Mergeable::Conflicting, + _ => Mergeable::Unknown, + }, + merge_state: p + .get("mergeStateStatus") + .and_then(Value::as_str) + .unwrap_or("UNKNOWN") + .to_string(), + review_decision: p + .get("reviewDecision") + .and_then(Value::as_str) + .map(str::to_string), + auto_merge_armed: p.get("autoMergeRequest").is_some_and(Value::is_object), + rollup: p + .pointer("/commits/nodes/0/commit/statusCheckRollup/state") + .and_then(Value::as_str) + .map(str::to_string), + updated_at: s("updatedAt"), + }); + } + RepoRead::Read { + gate, + truncated: total > prs.len(), + prs, + } +} + +/// Read gates and open PRs for `repos` in batches of `batch_size`. Returns +/// one read per repo, in order, and the number of queries spent. +/// +/// A batch that fails *as a whole* (a gateway 502, a timeout, a body with no +/// `data`) is split in half and each half retried, down to a single repo, +/// which gets one more attempt. A heavy repo therefore cannot take its +/// neighbours down with it. Per-repo errors inside a good response are not +/// retried: they are about that repo. +/// +/// When the remaining GraphQL budget would not cover another query of the +/// same cost, every later repo is [`RepoRead::Unavailable`] with the reset +/// time — never skipped silently. +pub fn fetch_board( + transport: &dyn GraphQlTransport, + repos: &[RepoId], + batch_size: usize, +) -> (Vec, usize) { + let mut st = FetchState::default(); + let mut out: Vec> = vec![None; repos.len()]; + let valid: Vec = (0..repos.len()) + .filter(|&i| { + let ok = split_slug(&repos[i]).is_some(); + if !ok { + out[i] = Some(RepoRead::Unavailable("not an `owner/name` slug".into())); + } + ok + }) + .collect(); + for chunk in valid.chunks(batch_size.max(1)) { + let batch: Vec = chunk.iter().map(|&i| repos[i].clone()).collect(); + for (&i, read) in chunk.iter().zip(read_chunk(transport, &batch, &mut st, 1)) { + out[i] = Some(read); + } + } + let reads = out + .into_iter() + .map(|r| r.unwrap_or_else(|| RepoRead::Unavailable("not fetched".into()))) + .collect(); + (reads, st.queries) +} + +#[derive(Default)] +struct FetchState { + queries: usize, + exhausted: Option, +} + +/// One batch, split-and-retried on whole-batch failure (see [`fetch_board`]). +fn read_chunk( + transport: &dyn GraphQlTransport, + batch: &[RepoId], + st: &mut FetchState, + retries: u8, +) -> Vec { + if let Some(reason) = &st.exhausted { + return batch + .iter() + .map(|_| RepoRead::Unavailable(reason.clone())) + .collect(); + } + let (doc, vars) = board_query(batch); + st.queries += 1; + let failure = match transport.execute(&json!({ "query": doc, "variables": vars })) { + Ok(resp) if resp.get("data").is_some_and(Value::is_object) => { + let (reads, rate) = parse_board_response(batch, &resp); + if let Some(r) = rate { + if r.remaining < r.cost.max(1) { + st.exhausted = Some(format!( + "GraphQL rate budget exhausted ({} left){}", + r.remaining, + r.reset_at + .as_deref() + .map(|t| format!(", resets at {t}")) + .unwrap_or_default() + )); + } + } + return reads; + } + Ok(resp) => error_text(&resp), + Err(reason) => reason, + }; + if batch.len() > 1 { + let (a, b) = batch.split_at(batch.len() / 2); + let mut reads = read_chunk(transport, a, st, retries); + reads.extend(read_chunk(transport, b, st, retries)); + reads + } else if retries > 0 { + read_chunk(transport, batch, st, retries - 1) + } else { + vec![RepoRead::Unavailable(failure)] + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::cell::RefCell; + + const SCHEMA: &str = include_str!("../graphql/github-schema.graphql"); + + fn validate(doc: &str) -> Result<(), String> { + use apollo_compiler::{ExecutableDocument, Schema}; + let schema = Schema::parse_and_validate(SCHEMA, "github-schema.graphql") + .map_err(|e| format!("schema: {}", e.errors))?; + ExecutableDocument::parse_and_validate(&schema, doc, "query.graphql") + .map(|_| ()) + .map_err(|e| e.errors.to_string()) + } + + #[test] + fn estate_repos_validates() { + validate(ESTATE_REPOS).unwrap(); + } + + #[test] + fn board_query_validates_at_every_batch_size() { + for n in [1, 2, BOARD_BATCH] { + let batch: Vec = (0..n).map(|i| RepoId::new(format!("o{i}/r{i}"))).collect(); + let (doc, vars) = board_query(&batch); + validate(&doc).unwrap_or_else(|e| panic!("batch {n}: {e}\n{doc}")); + assert_eq!(vars.as_object().unwrap().len(), 2 * n); + } + } + + #[test] + fn validator_rejects_a_made_up_field_in_the_board_fragment() { + let bad = BOARD_REPO.replace("mergeStateStatus", "mergeStateStatusX"); + let (doc, _) = board_query(&[RepoId::new("o/r")]); + let doc = doc.replace(BOARD_REPO, &bad); + assert!(validate(&doc).is_err()); + } + + /// Replays canned responses in order and records every request body. + struct Replay { + responses: RefCell>, + seen: RefCell>, + } + + impl Replay { + fn new(rs: Vec) -> Self { + Self { + responses: RefCell::new(rs.into_iter().rev().collect()), + seen: RefCell::new(Vec::new()), + } + } + } + + impl GraphQlTransport for Replay { + fn execute(&self, body: &Value) -> Result { + self.seen.borrow_mut().push(body.clone()); + self.responses + .borrow_mut() + .pop() + .ok_or_else(|| "no more canned responses".into()) + } + } + + fn page(total: u64, names: &[(&str, u64)], next: Option<&str>) -> Value { + json!({ "data": { + "rateLimit": { "cost": 1, "remaining": 4000, "resetAt": "2026-10-01T15:00:00Z" }, + "repositoryOwner": { "repositories": { + "totalCount": total, + "pageInfo": { "hasNextPage": next.is_some(), "endCursor": next }, + "nodes": names.iter().map(|(n, c)| json!({ + "nameWithOwner": n, "pullRequests": { "totalCount": c } + })).collect::>() + }} + }}) + } + + #[test] + fn enumeration_follows_cursors_and_sums_prs() { + let t = Replay::new(vec![ + page(3, &[("me/a", 2), ("me/b", 0)], Some("C1")), + page(3, &[("me/c", 5)], None), + ]); + let l = enumerate_owner(&t, "me").unwrap(); + assert_eq!(l.repos.len(), 3); + assert_eq!(l.open_prs(), 7); + assert_eq!(t.seen.borrow()[1]["variables"]["after"], "C1"); + } + + #[test] + fn a_transient_listing_failure_is_retried() { + struct Flaky(RefCell); + impl GraphQlTransport for Flaky { + fn execute(&self, _: &Value) -> Result { + *self.0.borrow_mut() += 1; + if *self.0.borrow() == 1 { + return Err("gh: HTTP 502".into()); + } + Ok(page(1, &[("me/a", 1)], None)) + } + } + let l = enumerate_owner(&Flaky(RefCell::new(0)), "me").unwrap(); + assert_eq!(l.repos.len(), 1); + } + + #[test] + fn enumeration_short_of_total_count_is_an_error() { + let t = Replay::new(vec![page(5, &[("me/a", 1)], None)]); + let e = enumerate_owner(&t, "me").unwrap_err(); + assert!(e.contains("reports 5 repositories but 1"), "{e}"); + } + + #[test] + fn enumeration_without_data_is_an_error_not_an_empty_estate() { + let t = Replay::new(vec![json!({ "errors": [{ "message": "rate limited" }] })]); + let e = enumerate_owner(&t, "me").unwrap_err(); + assert!(e.contains("rate limited"), "{e}"); + } + + fn pr_node(n: u64) -> Value { + json!({ + "number": n, "title": "t", "url": format!("https://github.com/me/a/pull/{n}"), + "isDraft": false, "mergeable": "MERGEABLE", "mergeStateStatus": "BLOCKED", + "reviewDecision": null, "updatedAt": "2026-10-01T10:00:00Z", + "headRefName": "h", "author": { "login": "bot" }, + "autoMergeRequest": null, + "commits": { "nodes": [ { "commit": { "statusCheckRollup": { "state": "FAILURE" } } } ] } + }) + } + + #[test] + fn gate_unions_rulesets_and_classic_protection_without_double_counting() { + let node = json!({ + "nameWithOwner": "me/a", + "defaultBranchRef": { + "name": "main", + "rules": { "nodes": [ + { "type": "REQUIRED_STATUS_CHECKS", "parameters": { + "requiredStatusChecks": [ { "context": "build" }, { "context": "test" } ] } }, + { "type": "PULL_REQUEST", "parameters": { "requiredApprovingReviewCount": 1 } }, + { "type": "DELETION", "parameters": null } + ]}, + "branchProtectionRule": { + "requiresStatusChecks": true, + "requiredStatusCheckContexts": ["test", "lint"], + "requiresApprovingReviews": true, + "requiredApprovingReviewCount": 2 + } + }, + "pullRequests": { "totalCount": 1, "nodes": [ pr_node(7) ] } + }); + let RepoRead::Read { + gate, + prs, + truncated, + } = parse_board_repo(&RepoId::new("me/a"), &node) + else { + panic!("expected a read") + }; + assert_eq!(gate.required_contexts, 3); // build, test, lint + assert_eq!(gate.required_approvals, 2); + assert!(!truncated); + assert_eq!(prs[0].rollup.as_deref(), Some("FAILURE")); + assert!(!prs[0].auto_merge_armed); + } + + #[test] + fn no_rules_and_no_protection_is_a_zero_gate() { + let node = json!({ + "nameWithOwner": "me/a", + "defaultBranchRef": { "name": "main", "rules": { "nodes": [] }, "branchProtectionRule": null }, + "pullRequests": { "totalCount": 150, "nodes": [ pr_node(1) ] } + }); + let RepoRead::Read { + gate, truncated, .. + } = parse_board_repo(&RepoId::new("me/a"), &node) + else { + panic!("expected a read") + }; + assert_eq!(gate.required_contexts, 0); + assert!(truncated, "150 open but 1 returned must be flagged"); + } + + #[test] + fn a_null_alias_is_unavailable_with_the_error_message() { + let resp = json!({ + "data": { "rateLimit": { "cost": 1, "remaining": 10, "resetAt": null }, "r0": null }, + "errors": [ { "path": ["r0"], "type": "NOT_FOUND", "message": "Could not resolve" } ] + }); + let (reads, _) = parse_board_response(&[RepoId::new("me/gone")], &resp); + assert_eq!(reads[0], RepoRead::Unavailable("Could not resolve".into())); + } + + /// Fails any query naming `heavy`; fails the first `flaky_first` calls + /// outright; otherwise answers every alias with an empty repo. + struct Gateway { + calls: RefCell, + flaky_first: usize, + } + + impl GraphQlTransport for Gateway { + fn execute(&self, body: &Value) -> Result { + *self.calls.borrow_mut() += 1; + if *self.calls.borrow() <= self.flaky_first { + return Err("gh: HTTP 502".into()); + } + let vars = body["variables"].as_object().unwrap(); + if vars.values().any(|v| v == "heavy") { + return Ok(json!({ "message": "timeout" })); + } + let mut data = serde_json::Map::new(); + data.insert( + "rateLimit".into(), + json!({ "cost": 1, "remaining": 4000, "resetAt": null }), + ); + for i in 0..vars.len() / 2 { + data.insert( + format!("r{i}"), + json!({ "nameWithOwner": "x", "defaultBranchRef": null, + "pullRequests": { "totalCount": 0, "nodes": [] } }), + ); + } + Ok(json!({ "data": data })) + } + } + + #[test] + fn one_heavy_repo_does_not_sink_its_batch() { + let t = Gateway { + calls: RefCell::new(0), + flaky_first: 0, + }; + let repos: Vec = ["me/a", "me/heavy", "me/b", "me/c"] + .iter() + .map(|s| RepoId::new(*s)) + .collect(); + let (reads, _) = fetch_board(&t, &repos, 4); + for i in [0, 2, 3] { + assert!( + matches!(reads[i], RepoRead::Read { .. }), + "repo {i}: {:?}", + reads[i] + ); + } + assert_eq!(reads[1], RepoRead::Unavailable("timeout".into())); + } + + #[test] + fn a_transient_gateway_failure_is_recovered() { + let t = Gateway { + calls: RefCell::new(0), + flaky_first: 1, + }; + let (reads, q) = fetch_board(&t, &[RepoId::new("me/a")], 4); + assert!(matches!(reads[0], RepoRead::Read { .. }), "{:?}", reads[0]); + assert_eq!(q, 2); + } + + #[test] + fn an_exhausted_budget_marks_the_rest_unavailable() { + let low = json!({ "data": { + "rateLimit": { "cost": 5, "remaining": 2, "resetAt": "2026-10-01T15:00:00Z" }, + "r0": { "nameWithOwner": "me/a", "defaultBranchRef": null, + "pullRequests": { "totalCount": 0, "nodes": [] } } + }}); + let t = Replay::new(vec![low]); + let repos = vec![RepoId::new("me/a"), RepoId::new("me/b")]; + let (reads, q) = fetch_board(&t, &repos, 1); + assert_eq!(q, 1); + assert!(matches!(reads[0], RepoRead::Read { .. })); + assert!(matches!(&reads[1], RepoRead::Unavailable(r) if r.contains("exhausted"))); + } +} diff --git a/crates/squabble-forge/src/inbox.rs b/crates/squabble-forge/src/inbox.rs new file mode 100644 index 0000000..3faef13 --- /dev/null +++ b/crates/squabble-forge/src/inbox.rs @@ -0,0 +1,209 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +//! Forge reads for `squabble inbox-sweep`: the state of every issue or pull +//! request a notification thread points at, batched through GraphQL. +//! +//! Each alias carries its own `number`, so the per-subject selection is +//! generated inline rather than from a shared fragment; it is still +//! schema-validated in tests. A subject that could not be read comes back +//! [`SubjectState::Unknown`] with the reason, and the sweep keeps its thread. + +use crate::{split_slug, GraphQlTransport}; +use serde_json::{json, Value}; +use squabble_core::inbox::{SubjectRef, SubjectState}; +use std::collections::BTreeMap; + +/// Subjects per state query. +pub const STATE_BATCH: usize = 50; + +/// Build one batched state query for `batch`, aliased `s0…sN`. +pub fn state_query(batch: &[SubjectRef]) -> (String, Value) { + let mut params = Vec::new(); + let mut fields = Vec::new(); + let mut vars = serde_json::Map::new(); + for (i, s) in batch.iter().enumerate() { + let (o, n) = split_slug(&s.repo).unwrap_or(("", "")); + params.push(format!("$o{i}: String!, $n{i}: String!, $k{i}: Int!")); + fields.push(format!( + " s{i}: repository(owner: $o{i}, name: $n{i}) {{ issueOrPullRequest(number: $k{i}) {{ __typename ... on PullRequest {{ prState: state }} ... on Issue {{ issueState: state }} }} }}" + )); + vars.insert(format!("o{i}"), json!(o)); + vars.insert(format!("n{i}"), json!(n)); + vars.insert(format!("k{i}"), json!(s.number)); + } + let doc = format!( + "query SubjectStates({}) {{\n rateLimit {{ cost remaining resetAt }}\n{}\n}}\n", + params.join(", "), + fields.join("\n") + ); + (doc, Value::Object(vars)) +} + +/// Read one alias of a state response. +fn state_of(resp: &Value, alias: &str) -> SubjectState { + let node = resp.pointer(&format!("/data/{alias}/issueOrPullRequest")); + let state = node.and_then(|n| n.get("prState").or_else(|| n.get("issueState"))); + match state.and_then(Value::as_str) { + Some("OPEN") => SubjectState::Open, + Some("MERGED") => SubjectState::Merged, + Some("CLOSED") => SubjectState::Closed, + Some(other) => SubjectState::Unknown(format!("unrecognised state `{other}`")), + None => { + let msg = resp + .get("errors") + .and_then(Value::as_array) + .into_iter() + .flatten() + .find(|e| e.pointer("/path/0").and_then(Value::as_str) == Some(alias)) + .and_then(|e| e.get("message").and_then(Value::as_str)) + .map(str::to_string); + SubjectState::Unknown(msg.unwrap_or_else(|| { + if resp.get("data").is_some_and(Value::is_object) { + "subject not found".into() + } else { + "response carried no data".into() + } + })) + } + } +} + +/// Look up the state of every subject. Duplicates are read once. Returns the +/// states and the number of queries spent. +pub fn fetch_states( + transport: &dyn GraphQlTransport, + subjects: &[SubjectRef], + batch_size: usize, +) -> (BTreeMap, usize) { + let mut unique: Vec = subjects.to_vec(); + unique.sort(); + unique.dedup(); + let mut out = BTreeMap::new(); + let (valid, invalid): (Vec, Vec) = unique + .into_iter() + .partition(|s| split_slug(&s.repo).is_some()); + for s in invalid { + out.insert(s, SubjectState::Unknown("not an `owner/name` slug".into())); + } + let mut queries = 0; + let mut exhausted: Option = None; + for chunk in valid.chunks(batch_size.max(1)) { + if let Some(r) = &exhausted { + for s in chunk { + out.insert(s.clone(), SubjectState::Unknown(r.clone())); + } + continue; + } + let (doc, vars) = state_query(chunk); + queries += 1; + match transport.execute(&json!({ "query": doc, "variables": vars })) { + Ok(resp) => { + for (i, s) in chunk.iter().enumerate() { + out.insert(s.clone(), state_of(&resp, &format!("s{i}"))); + } + let rem = resp + .pointer("/data/rateLimit/remaining") + .and_then(Value::as_u64); + let cost = resp.pointer("/data/rateLimit/cost").and_then(Value::as_u64); + if let (Some(rem), Some(cost)) = (rem, cost) { + if rem < cost.max(1) { + exhausted = Some(format!("GraphQL rate budget exhausted ({rem} left)")); + } + } + } + Err(e) => { + for s in chunk { + out.insert(s.clone(), SubjectState::Unknown(e.clone())); + } + } + } + } + (out, queries) +} + +#[cfg(test)] +mod tests { + use super::*; + use squabble_core::chains::RepoId; + + const SCHEMA: &str = include_str!("../graphql/github-schema.graphql"); + + fn validate(doc: &str) -> Result<(), String> { + use apollo_compiler::{ExecutableDocument, Schema}; + let schema = Schema::parse_and_validate(SCHEMA, "github-schema.graphql") + .map_err(|e| format!("schema: {}", e.errors))?; + ExecutableDocument::parse_and_validate(&schema, doc, "query.graphql") + .map(|_| ()) + .map_err(|e| e.errors.to_string()) + } + + fn subj(r: &str, n: u64) -> SubjectRef { + SubjectRef { + repo: RepoId::new(r), + number: n, + } + } + + #[test] + fn state_query_validates_at_every_batch_size() { + for n in [1, 3, STATE_BATCH] { + let batch: Vec = (0..n) + .map(|i| subj(&format!("o/r{i}"), i as u64 + 1)) + .collect(); + let (doc, vars) = state_query(&batch); + validate(&doc).unwrap_or_else(|e| panic!("batch {n}: {e}")); + assert_eq!(vars.as_object().unwrap().len(), 3 * n); + } + } + + #[test] + fn validator_rejects_a_wrong_field_in_the_generated_query() { + let (doc, _) = state_query(&[subj("o/r", 1)]); + assert!(validate(&doc.replace("prState: state", "prState: stateX")).is_err()); + } + + #[test] + fn states_parse_and_failures_stay_unknown() { + let resp = json!({ + "data": { + "s0": { "issueOrPullRequest": { "__typename": "PullRequest", "prState": "MERGED" } }, + "s1": { "issueOrPullRequest": { "__typename": "Issue", "issueState": "CLOSED" } }, + "s2": { "issueOrPullRequest": { "__typename": "PullRequest", "prState": "OPEN" } }, + "s3": null, + "s4": { "issueOrPullRequest": null } + }, + "errors": [ { "path": ["s3"], "message": "Could not resolve to a Repository" } ] + }); + assert_eq!(state_of(&resp, "s0"), SubjectState::Merged); + assert_eq!(state_of(&resp, "s1"), SubjectState::Closed); + assert_eq!(state_of(&resp, "s2"), SubjectState::Open); + assert_eq!( + state_of(&resp, "s3"), + SubjectState::Unknown("Could not resolve to a Repository".into()) + ); + assert_eq!( + state_of(&resp, "s4"), + SubjectState::Unknown("subject not found".into()) + ); + assert_eq!( + state_of(&json!({ "message": "502" }), "s0"), + SubjectState::Unknown("response carried no data".into()) + ); + } + + #[test] + fn duplicates_are_read_once_and_transport_errors_are_unknown() { + struct Down; + impl GraphQlTransport for Down { + fn execute(&self, _: &Value) -> Result { + Err("gh: HTTP 502".into()) + } + } + let (m, q) = fetch_states(&Down, &[subj("o/r", 1), subj("o/r", 1), subj("o/r", 2)], 50); + assert_eq!(q, 1); + assert_eq!(m.len(), 2); + assert!(m + .values() + .all(|s| matches!(s, SubjectState::Unknown(r) if r.contains("502")))); + } +} diff --git a/crates/squabble-forge/src/lib.rs b/crates/squabble-forge/src/lib.rs index 0e835b9..06255d6 100644 --- a/crates/squabble-forge/src/lib.rs +++ b/crates/squabble-forge/src/lib.rs @@ -26,6 +26,9 @@ use squabble_core::chains::{RepoId, RepoSnapshot, ScanStatus, SourceCost, Workfl use std::io::Write; use std::process::{Command, Stdio}; +pub mod board; +pub mod inbox; + /// The per-repo fragment, aliased once per repo in [`chains_query`]. pub const WORKFLOW_TREE: &str = include_str!("../graphql/workflow_tree.graphql");