From 5677e58fe178edb594ed4f515baf862ca63adce9 Mon Sep 17 00:00:00 2001 From: mintaka Date: Wed, 7 Oct 2026 00:00:48 -0400 Subject: [PATCH] feat(forge): read PR timestamps, closing refs and PR list rows (RIG-4034) forge.PullRequest gains CreatedAt, UpdatedAt and ClosingRefs. The full GraphQL read fetches closingIssuesReferences on its first page only. ListUpdatedIssues returns the PR rows beside the issues, and ListOpenPullRequests lists open PRs for the backfill. The store gains PullRequestUpdatedAt, PRsBackfilledAt and MarkPRsBackfilled. Spec-impact: none. Refs RIG-4034 Co-authored-by: Matt Wilkinson --- go/internal/forge/github.go | 44 ++++--- go/internal/forge/github_graphql.go | 98 ++++++++++---- go/internal/forge/github_graphql_test.go | 7 +- go/internal/forge/github_test.go | 2 +- go/internal/forge/golden_capture_test.go | 9 +- go/internal/forge/notify_reader.go | 53 ++++++-- go/internal/forge/notify_reader_test.go | 31 +++-- go/internal/forge/provider.go | 11 ++ go/internal/forge/pull_link_reads_test.go | 123 ++++++++++++++++++ .../testdata/github/create_pull_request.json | 5 +- .../testdata/github/get_pull_request.json | 8 +- .../github/transition_pull_request_close.json | 5 +- .../transition_pull_request_reopen.json | 5 +- go/internal/ingest/board_reconcile.go | 4 +- go/internal/ingest/board_reconcile_test.go | 25 ++-- go/internal/store/db/forge_cursors.sql.go | 43 ++++++ go/internal/store/db/pull_requests.sql.go | 24 ++++ go/internal/store/db/querier.go | 3 + go/internal/store/forge_cursors.go | 44 +++++++ go/internal/store/pull_requests.go | 20 +++ .../store/pull_requests_pgtest_test.go | 34 +++++ go/internal/store/queries/forge_cursors.sql | 8 ++ go/internal/store/queries/pull_requests.sql | 4 + 23 files changed, 528 insertions(+), 82 deletions(-) create mode 100644 go/internal/forge/pull_link_reads_test.go diff --git a/go/internal/forge/github.go b/go/internal/forge/github.go index 1a0e95e8a..79cb9a8ea 100644 --- a/go/internal/forge/github.go +++ b/go/internal/forge/github.go @@ -309,13 +309,15 @@ func (r ghComment) toComment() Comment { // the fields forge.PullRequest needs at create time are decoded; the read-side // roll-ups (Changed/Checks/Reviews/Threads) are GetPullRequest's (ghPullDetail). type ghPull struct { - Number uint64 `json:"number"` - Title string `json:"title"` - Body string `json:"body"` - State string `json:"state"` - HTMLURL string `json:"html_url"` - Draft bool `json:"draft"` - Head struct { + Number uint64 `json:"number"` + Title string `json:"title"` + Body string `json:"body"` + State string `json:"state"` + HTMLURL string `json:"html_url"` + Draft bool `json:"draft"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` + Head struct { Ref string `json:"ref"` } `json:"head"` Base struct { @@ -339,9 +341,20 @@ func (r ghPull) toPullRequest() PullRequest { BaseRef: r.Base.Ref, ForgeAccount: r.User.Login, Draft: r.Draft, + CreatedAt: parseGHTime(r.CreatedAt), + UpdatedAt: parseGHTime(r.UpdatedAt), } } +// parseGHTime parses a GitHub RFC-3339 timestamp; empty or malformed is the zero time. +func parseGHTime(s string) time.Time { + t, err := time.Parse(time.RFC3339, s) + if err != nil { + return time.Time{} + } + return t +} + // CreateIssue creates an issue on repo. in.Body is PRE-stamped by the Service // (DL-050); the Provider sends it verbatim. Returns the created forge.Issue. func (g *GitHub) CreateIssue(ctx context.Context, repo string, in CreateIssue) (Issue, error) { @@ -564,6 +577,8 @@ type ghPullDetail struct { Deletions uint32 `json:"deletions"` ChangedFiles uint32 `json:"changed_files"` Merged bool `json:"merged"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` Head struct { Ref string `json:"ref"` SHA string `json:"sha"` @@ -600,6 +615,8 @@ func (r ghPullDetail) toPullRequest() PullRequest { Additions: r.Additions, Deletions: r.Deletions, }, + CreatedAt: parseGHTime(r.CreatedAt), + UpdatedAt: parseGHTime(r.UpdatedAt), } } @@ -646,12 +663,13 @@ func (g *GitHub) GetPullRequest(ctx context.Context, repo string, number uint64) }) } - checks, threads, err := g.checksForPull(ctx, coord, detail.Head.SHA, true) + checks, gql, err := g.checksForPull(ctx, coord, detail.Head.SHA, true) if err != nil { return PullRequest{}, fmt.Errorf("forge: github get pull request %q#%d: %w", repo, number, err) } pr.Checks = checks - pr.Threads = threads + pr.Threads = gql.threads + pr.ClosingRefs = gql.closingRefs return pr, nil } @@ -1021,12 +1039,6 @@ func (r ghIssue) toIssue() Issue { for _, l := range r.Labels { labels = append(labels, l.Name) } - var updated time.Time - if r.UpdatedAt != "" { - if t, err := time.Parse(time.RFC3339, r.UpdatedAt); err == nil { - updated = t - } - } return Issue{ Number: r.Number, Title: r.Title, @@ -1035,7 +1047,7 @@ func (r ghIssue) toIssue() Issue { URL: r.HTMLURL, ForgeAccount: r.User.Login, Labels: labels, - UpdatedAt: updated, + UpdatedAt: parseGHTime(r.UpdatedAt), } } diff --git a/go/internal/forge/github_graphql.go b/go/internal/forge/github_graphql.go index 4e6c8b252..d6ad5854d 100644 --- a/go/internal/forge/github_graphql.go +++ b/go/internal/forge/github_graphql.go @@ -9,6 +9,7 @@ import ( "context" "errors" "fmt" + "log/slog" "math" "net/http" "strings" @@ -172,12 +173,29 @@ type ghGQLCommitNode struct { } `json:"commit"` } +// ghGQLClosingRef is one issue the PR closes. +type ghGQLClosingRef struct { + Number uint64 `json:"number"` + Repository struct { + NameWithOwner string `json:"nameWithOwner"` + } `json:"repository"` +} + +// ghGQLClosingRefs is the issues the PR closes (closingIssuesReferences). +type ghGQLClosingRefs struct { + Nodes []ghGQLClosingRef `json:"nodes"` +} + // ghGQLPull is the pull request half; each connection is absent when excluded. type ghGQLPull struct { ReviewThreads *ghGQLThreads `json:"reviewThreads"` Commits *ghGQLPullCommits `json:"commits"` + ClosingRefs *ghGQLClosingRefs `json:"closingIssuesReferences"` } +// closingRefsCap is the closingIssuesReferences page size; a PR closing more is truncated. +const closingRefsCap = 25 + // ghGQLRepo is the repository root of pullReadQuery. type ghGQLRepo struct { PullRequest *ghGQLPull `json:"pullRequest"` @@ -199,9 +217,12 @@ type ghGQLThreadNode struct { // @include flags let a later page fetch only the connection still paging. The // contexts are read through the pull request (commits(last: 1)), so only Pull // requests: read is needed; the commit oid pins each page to the REST head SHA. -const pullReadQuery = `query($owner: String!, $name: String!, $number: Int!, $threads: Boolean!, $threadsAfter: String, $contexts: Boolean!, $contextsAfter: String) { +const pullReadQuery = `query($owner: String!, $name: String!, $number: Int!, $threads: Boolean!, $threadsAfter: String, $contexts: Boolean!, $contextsAfter: String, $refs: Boolean!) { repository(owner: $owner, name: $name) { pullRequest(number: $number) { + closingIssuesReferences(first: 25) @include(if: $refs) { + nodes { number repository { nameWithOwner } } + } reviewThreads(first: 100, after: $threadsAfter) @include(if: $threads) { pageInfo { hasNextPage endCursor } nodes { @@ -256,25 +277,31 @@ func (g *GitHub) graphQLURL() string { return "https://" + g.host + "/api/graphql" } +// pullReadExtras is what a full GraphQL read adds beyond checks. +type pullReadExtras struct { + threads []ReviewThread + closingRefs []IssueRef +} + // checksForPull is the one checks path for a pull request (GetPullRequest and // Checks): the REST roll-up from checksForSHA, with Required set from the PR's -// required contexts. withThreads also returns the review threads from the same -// GraphQL leg, so GetPullRequest pays one leg, not two. -func (g *GitHub) checksForPull(ctx context.Context, c pullCoord, sha string, withThreads bool) (Checks, []ReviewThread, error) { +// required contexts. full also returns the review threads and closing references +// from the same GraphQL leg, so GetPullRequest pays one leg, not two. +func (g *GitHub) checksForPull(ctx context.Context, c pullCoord, sha string, full bool) (Checks, pullReadExtras, error) { checks, err := g.checksForSHA(ctx, c.repo, sha) if err != nil { - return Checks{}, nil, err + return Checks{}, pullReadExtras{}, err } - threads, required, err := g.pullGraphQL(ctx, c, sha, withThreads) + w, err := g.pullGraphQL(ctx, c, sha, full) if err != nil { - return Checks{}, nil, fmt.Errorf("forge: github graphql for %q#%d: %w", c.repo, c.number, err) + return Checks{}, pullReadExtras{}, fmt.Errorf("forge: github graphql for %q#%d: %w", c.repo, c.number, err) } for i := range checks.Checks { - if _, ok := required[checks.Checks[i].Name]; ok { + if _, ok := w.required[checks.Checks[i].Name]; ok { checks.Checks[i].Required = true } } - return checks, threads, nil + return checks, pullReadExtras{threads: w.threads, closingRefs: w.closingRefs}, nil } // pullGraphQLWalk is the cursor state of one pullGraphQL walk. A connection @@ -282,48 +309,71 @@ func (g *GitHub) checksForPull(ctx context.Context, c pullCoord, sha string, wit type pullGraphQLWalk struct { sha string // REST head SHA every contexts page must match threads []ReviewThread + closingRefs []IssueRef required map[string]struct{} threadsAfter, contextsAfter any // nil sends JSON null: the first page moreThreads, moreContexts bool + refs bool // read closing references; first page only } -// pullGraphQL walks pullReadQuery to completion: every review-thread page (when -// withThreads) and every context page of head commit sha. It returns the threads -// in forge order and the set of required context names (a CheckRun's name, a -// StatusContext's context), which match the REST check names. -func (g *GitHub) pullGraphQL(ctx context.Context, c pullCoord, sha string, withThreads bool) ([]ReviewThread, map[string]struct{}, error) { - w := pullGraphQLWalk{sha: sha, required: map[string]struct{}{}, moreThreads: withThreads, moreContexts: true} +// pullGraphQL walks pullReadQuery to completion: every review-thread page and +// the closing references (when full) and every context page of head commit sha. +// The walk holds the threads in forge order and the set of required context +// names (a CheckRun's name, a StatusContext's context), which match REST's. +func (g *GitHub) pullGraphQL(ctx context.Context, c pullCoord, sha string, full bool) (*pullGraphQLWalk, error) { + w := &pullGraphQLWalk{sha: sha, required: map[string]struct{}{}, moreThreads: full, moreContexts: true, refs: full} for w.moreThreads || w.moreContexts { if err := ctx.Err(); err != nil { - return nil, nil, err + return nil, err } data, err := graphQL[ghGQLPullData](ctx, g, pullReadQuery, map[string]any{ "owner": c.owner, "name": c.name, "number": c.number, "threads": w.moreThreads, "threadsAfter": w.threadsAfter, "contexts": w.moreContexts, "contextsAfter": w.contextsAfter, + "refs": w.refs, }) if err != nil { - return nil, nil, err + return nil, err } if data.Repository == nil { - return nil, nil, fmt.Errorf("forge: github graphql: repository %q not found", c.repo) + return nil, fmt.Errorf("forge: github graphql: repository %q not found", c.repo) } pr := data.Repository.PullRequest if pr == nil { - return nil, nil, fmt.Errorf("forge: github graphql: pull request %q#%d not found", c.repo, c.number) + return nil, fmt.Errorf("forge: github graphql: pull request %q#%d not found", c.repo, c.number) + } + if w.refs { + if err := foldClosingRefs(ctx, w, c, pr.ClosingRefs); err != nil { + return nil, err + } } if w.moreThreads { - if err := g.foldThreads(ctx, &w, c, pr.ReviewThreads); err != nil { - return nil, nil, err + if err := g.foldThreads(ctx, w, c, pr.ReviewThreads); err != nil { + return nil, err } } if w.moreContexts { - if err := foldContexts(&w, c, pr.Commits); err != nil { - return nil, nil, err + if err := foldContexts(w, c, pr.Commits); err != nil { + return nil, err } } } - return w.threads, w.required, nil + return w, nil +} + +// foldClosingRefs records the PR's closing references and turns them off for later pages. +func foldClosingRefs(ctx context.Context, w *pullGraphQLWalk, c pullCoord, conn *ghGQLClosingRefs) error { + if conn == nil { + return fmt.Errorf("forge: github graphql: pull request %q#%d has no closing references connection", c.repo, c.number) + } + w.refs = false + for _, n := range conn.Nodes { + w.closingRefs = append(w.closingRefs, IssueRef{Repo: n.Repository.NameWithOwner, Number: n.Number}) + } + if len(conn.Nodes) >= closingRefsCap { + slog.WarnContext(ctx, "forge: github closing references truncated", "repo", c.repo, "number", c.number, "cap", closingRefsCap) + } + return nil } // foldThreads appends one page of review threads to w and advances its cursor. diff --git a/go/internal/forge/github_graphql_test.go b/go/internal/forge/github_graphql_test.go index 493a8e754..885e62923 100644 --- a/go/internal/forge/github_graphql_test.go +++ b/go/internal/forge/github_graphql_test.go @@ -32,6 +32,7 @@ func noRollupGraphQL(sha string) string { // status-check rollup on the last commit, head sha. func emptyPullGraphQL(sha string) string { return `{"data":{"repository":{"pullRequest":{ + "closingIssuesReferences":{"nodes":[]}, "reviewThreads":{"pageInfo":{"hasNextPage":false,"endCursor":null},"nodes":[]}, "commits":{"nodes":[{"commit":{"oid":"` + sha + `","statusCheckRollup":null}}]}}}}}` } @@ -48,9 +49,10 @@ func gqlBody(t *testing.T, data any) string { // pullPage builds one pullReadQuery page. A nil threads or contexts leaves that // connection out, as @include(if: false) does; a nil rollup is a commit with no checks. +// Closing refs ride every page here; the walk reads them from the first only. func pullPage(t *testing.T, threads *ghGQLThreads, commits *ghGQLPullCommits) string { t.Helper() - return gqlBody(t, ghGQLPullData{Repository: &ghGQLRepo{PullRequest: &ghGQLPull{ReviewThreads: threads, Commits: commits}}}) + return gqlBody(t, ghGQLPullData{Repository: &ghGQLRepo{PullRequest: &ghGQLPull{ReviewThreads: threads, Commits: commits, ClosingRefs: &ghGQLClosingRefs{}}}}) } // threadsConn is one page of review threads. @@ -293,6 +295,9 @@ func TestPullGraphQLErrorBranches(t *testing.T) { {"null pull request", []scriptedResponse{ok200(`{"data":{"repository":{"pullRequest":null}}}`)}, "pull request"}, {"no threads connection", []scriptedResponse{ok200(pullPage(t, nil, rollup(headSHA, false, "c")))}, "review threads connection"}, {"no commits connection", []scriptedResponse{ok200(pullPage(t, threadsConn(false, "t"), nil))}, "commits connection"}, + {"no closing references connection", []scriptedResponse{ok200(gqlBody(t, ghGQLPullData{Repository: &ghGQLRepo{PullRequest: &ghGQLPull{ + ReviewThreads: threadsConn(false, "t"), Commits: rollup(headSHA, false, "c"), + }}}))}, "closing references connection"}, {"next page without a cursor", []scriptedResponse{ok200(pullPage(t, threadsConn(true, ""), rollup(headSHA, false, "c")))}, "without an endCursor"}, {"repeated cursor", []scriptedResponse{ ok200(pullPage(t, threadsConn(true, "SAME"), rollup(headSHA, false, "c"))), diff --git a/go/internal/forge/github_test.go b/go/internal/forge/github_test.go index 0e11499e4..849dba522 100644 --- a/go/internal/forge/github_test.go +++ b/go/internal/forge/github_test.go @@ -1425,7 +1425,7 @@ func TestGetIssueBudgetGateFailFast(t *testing.T) { // with a human and a bot comment, an unresolved PR-level thread, and one // required check run beside a non-required legacy status. const happyPullGraphQL = `{"data":{"repository":{ - "pullRequest":{"reviewThreads":{ + "pullRequest":{"closingIssuesReferences":{"nodes":[]},"reviewThreads":{ "pageInfo":{"hasNextPage":false,"endCursor":"t1"}, "nodes":[ {"id":"T1","isResolved":true,"path":"main.go","comments":{ diff --git a/go/internal/forge/golden_capture_test.go b/go/internal/forge/golden_capture_test.go index 268f513ce..512e396fc 100644 --- a/go/internal/forge/golden_capture_test.go +++ b/go/internal/forge/golden_capture_test.go @@ -33,6 +33,7 @@ const ( canonURL = "https://example.invalid/canonical" // html_url, url, target_url -> URL canonIDString = "canonical-id" // a string/UUID id (Linear) -> ID / resolve coordinate canonUpdatedAt = "2026-08-01T12:30:00Z" // updated_at, updatedAt -> UpdatedAt + canonCreatedAt = "2026-08-01T12:00:00Z" // created_at -> CreatedAt canonSHA = "canonicalsha" // sha -> HeadSHA canonRef = "canonical-ref" // ref -> HeadRef/BaseRef canonTitle = "canonical title" // title -> Title @@ -71,6 +72,7 @@ var volatileFields = map[string]struct{}{ "ID": {}, "URL": {}, "UpdatedAt": {}, + "CreatedAt": {}, "HeadSHA": {}, "HeadRef": {}, "BaseRef": {}, @@ -94,6 +96,7 @@ var wireVolatile = map[string]func(node any) any{ "target_url": fixedSentinel(canonURL), "updated_at": fixedSentinel(canonUpdatedAt), "updatedAt": fixedSentinel(canonUpdatedAt), + "created_at": fixedSentinel(canonCreatedAt), "login": fixedSentinel(canonAccount), "displayName": fixedSentinel(canonAccount), "sha": fixedSentinel(canonSHA), @@ -134,6 +137,7 @@ func canonCursor(node any) any { // - ID <- id (ghComment.ID numeric; Linear comment id is a UUID string) // - URL <- html_url,url,target_url (GitHub HTMLURL + ghStatus.TargetURL -> Check.URL; Linear URL) // - UpdatedAt<- updated_at,updatedAt (GitHub updated_at; Linear updatedAt) +// - CreatedAt<- created_at (ghPull/ghPullDetail.CreatedAt) // - HeadSHA <- sha (ghPullDetail.Head.SHA) // - HeadRef <- ref (ghPull/ghPullDetail.Head.Ref) // - BaseRef <- ref (ghPull/ghPullDetail.Base.Ref) — same wire key as HeadRef @@ -147,6 +151,7 @@ var domainToWire = map[string][]string{ "ID": {"id"}, "URL": {"html_url", "url", "target_url"}, "UpdatedAt": {"updated_at", "updatedAt"}, + "CreatedAt": {"created_at"}, "HeadSHA": {"sha"}, "HeadRef": {"ref"}, "BaseRef": {"ref"}, @@ -426,7 +431,7 @@ func TestUpdateCanonicalizeStable(t *testing.T) { // (1) Completeness across every wire-volatile key. all := json.RawMessage(`{ "number": 1, "id": 2, "html_url": "h", "url": "u", "target_url": "t", - "updated_at": "a", "updatedAt": "b", "login": "l", "displayName": "d", + "updated_at": "a", "updatedAt": "b", "created_at": "c", "login": "l", "displayName": "d", "sha": "s", "oid": "o", "ref": "r", "title": "ti", "body": "bo", "description": "de", "endCursor": "Y3Vyc29y", "last": { "endCursor": null }, "state": "open", "keep": "kept" @@ -435,6 +440,7 @@ func TestUpdateCanonicalizeStable(t *testing.T) { "number": 42, "id": 42, "html_url": "https://example.invalid/canonical", "url": "https://example.invalid/canonical", "target_url": "https://example.invalid/canonical", "updated_at": "2026-08-01T12:30:00Z", "updatedAt": "2026-08-01T12:30:00Z", + "created_at": "2026-08-01T12:00:00Z", "login": "octocat", "displayName": "octocat", "sha": "canonicalsha", "oid": "canonicalsha", "ref": "canonical-ref", "title": "canonical title", "body": "canonical body", "description": "canonical body", "endCursor": "canonical-cursor", "last": { "endCursor": null }, @@ -575,6 +581,7 @@ func TestUpdateCanonicalizeComposite(t *testing.T) { // extra leg 4: GraphQL threads + required contexts — a bot comment login // and body inside nested arrays, a thread node id, and paging cursors. {status: 200, body: json.RawMessage(`{ "data": { "repository": { "pullRequest": { + "closingIssuesReferences": { "nodes": [] }, "reviewThreads": { "pageInfo": { "hasNextPage": false, "endCursor": "live-cursor" }, "nodes": [ { "id": "PRRT_live", "isResolved": true, "path": "main.go", "comments": { "pageInfo": { "hasNextPage": false, "endCursor": "live-c" }, diff --git a/go/internal/forge/notify_reader.go b/go/internal/forge/notify_reader.go index 967c70cfc..9b30306bc 100644 --- a/go/internal/forge/notify_reader.go +++ b/go/internal/forge/notify_reader.go @@ -394,10 +394,23 @@ func ghWalkNewArtifacts[R any](ctx context.Context, g *GitHub, base string, sinc return ConditionalResult[[]Issue]{V: out, ETag: pageETag}, nil } +// UpdatedRows is one updated-order walk, split by kind. +type UpdatedRows struct { + Issues []Issue + Pulls []UpdatedPull +} + +// UpdatedPull is a PR row from the issues list: enough to decide whether to hydrate it. +type UpdatedPull struct { + Number uint64 + State string + UpdatedAt time.Time +} + // ListUpdatedIssues walks /repos/{repo}/issues?state=all&sort=updated& // direction=desc newest-updated-first (page 1 conditioned on etag; a 304 => -// NotModified), collecting issue rows (PR rows dropped by the pull_request -// marker, mirroring ListNewArtifacts) until a page's oldest updated_at is +// NotModified), collecting issue rows and, separately, the PR rows GitHub +// interleaves (told apart by the pull_request marker) until a page's oldest updated_at is // strictly < since, or no rel="next" remains. Rows with updated_at == since are // RE-included: GitHub's updated_at is second-granularity, so a <= stop would // permanently exclude an issue updated in the same second as the stored @@ -406,9 +419,9 @@ func ghWalkNewArtifacts[R any](ctx context.Context, g *GitHub, base string, sinc // ETag to re-store. This is the updated-order sibling of ListNewArtifacts // (created-order, number-keyed): the reconcile/backfill read cannot see updates // to existing issues via the created-order walk. -func (g *GitHub) ListUpdatedIssues(ctx context.Context, repo string, since time.Time, etag string) (ConditionalResult[[]Issue], error) { +func (g *GitHub) ListUpdatedIssues(ctx context.Context, repo string, since time.Time, etag string) (ConditionalResult[UpdatedRows], error) { base := g.apiBase() + "/repos/" + repo + "/issues?state=all&sort=updated&direction=desc" - var out []Issue + var out UpdatedRows pageETag := "" for page := 1; ; page++ { u := base + "&per_page=" + strconv.Itoa(perPage) + "&page=" + strconv.Itoa(page) @@ -419,11 +432,11 @@ func (g *GitHub) ListUpdatedIssues(ctx context.Context, repo string, since time. var rows []ghIssue notMod, e, hasNext, err := g.getJSONCond(ctx, u, sendETag, &rows) if err != nil { - return ConditionalResult[[]Issue]{}, fmt.Errorf("forge: github list updated issues %q: %w", repo, err) + return ConditionalResult[UpdatedRows]{}, fmt.Errorf("forge: github list updated issues %q: %w", repo, err) } if page == 1 { if notMod { - return ConditionalResult[[]Issue]{NotModified: true}, nil + return ConditionalResult[UpdatedRows]{NotModified: true}, nil } pageETag = e } @@ -441,18 +454,38 @@ func (g *GitHub) ListUpdatedIssues(ctx context.Context, repo string, since time. if iss.UpdatedAt.IsZero() { continue } - // Drop the PR rows GitHub interleaves into /issues (mirroring - // ListNewArtifacts / ListIssuesPage's pull_request-marker guard). if raw := r.PullRequest; raw != nil && len(*raw) > 0 && string(*raw) != jsonNull { + out.Pulls = append(out.Pulls, UpdatedPull{Number: iss.Number, State: iss.State, UpdatedAt: iss.UpdatedAt}) continue } - out = append(out, iss) + out.Issues = append(out.Issues, iss) } if reachedOld || !hasNext { break } } - return ConditionalResult[[]Issue]{V: out, ETag: pageETag}, nil + return ConditionalResult[UpdatedRows]{V: out, ETag: pageETag}, nil +} + +// ghOpenPull is one row of the open pulls list. +type ghOpenPull struct { + Number uint64 `json:"number"` + State string `json:"state"` + UpdatedAt string `json:"updated_at"` +} + +// ListOpenPullRequests lists every open PR on repo, for the backfill pass. +func (g *GitHub) ListOpenPullRequests(ctx context.Context, repo string) ([]UpdatedPull, error) { + rows, err := getAllPages(ctx, g, g.apiBase()+"/repos/"+repo+"/pulls?state=open", + func(e []ghOpenPull) []ghOpenPull { return e }) + if err != nil { + return nil, fmt.Errorf("forge: github list open pull requests %q: %w", repo, err) + } + out := make([]UpdatedPull, 0, len(rows)) + for _, r := range rows { + out = append(out, UpdatedPull{Number: r.Number, State: r.State, UpdatedAt: parseGHTime(r.UpdatedAt)}) + } + return out, nil } // --- Linear arm -------------------------------------------------------------- diff --git a/go/internal/forge/notify_reader_test.go b/go/internal/forge/notify_reader_test.go index f2858cabe..4fde3fe40 100644 --- a/go/internal/forge/notify_reader_test.go +++ b/go/internal/forge/notify_reader_test.go @@ -9,6 +9,7 @@ import ( "context" "errors" "net/http" + "slices" "testing" "time" @@ -305,11 +306,11 @@ func TestListUpdatedIssuesStopsStrictlyBelowSince(t *testing.T) { } // 50,49 (page 1) + 48 (== since, re-included) = 3; 40 is strictly below and // stops the walk before any page 3 is requested. - if len(res.V) != 3 { - t.Fatalf("kept %d (%+v), want 3 (50,49,48; 48 re-included at == since, 40 stops the walk)", len(res.V), res.V) + if len(res.V.Issues) != 3 { + t.Fatalf("kept %d (%+v), want 3 (50,49,48; 48 re-included at == since, 40 stops the walk)", len(res.V.Issues), res.V.Issues) } - if res.V[2].Number != 48 { - t.Errorf("V[2].Number = %d, want 48 (the == since row is re-included)", res.V[2].Number) + if res.V.Issues[2].Number != 48 { + t.Errorf("V[2].Number = %d, want 48 (the == since row is re-included)", res.V.Issues[2].Number) } if rt.calls != 2 { t.Errorf("calls = %d, want 2 (the walk stopped once a page's oldest was strictly < since)", rt.calls) @@ -342,11 +343,11 @@ func TestListUpdatedIssuesMalformedRowDoesNotTruncate(t *testing.T) { t.Fatalf("ListUpdatedIssues: %v", err) } // 50 and 48 are collected; the malformed 49 is skipped, not a stop signal. - if len(res.V) != 2 { - t.Fatalf("kept %d (%+v), want 2 (50,48; the malformed 49 is skipped without truncating)", len(res.V), res.V) + if len(res.V.Issues) != 2 { + t.Fatalf("kept %d (%+v), want 2 (50,48; the malformed 49 is skipped without truncating)", len(res.V.Issues), res.V.Issues) } - if res.V[0].Number != 50 || res.V[1].Number != 48 { - t.Errorf("kept %d,%d, want 50,48 (the row behind the malformed one still collected)", res.V[0].Number, res.V[1].Number) + if res.V.Issues[0].Number != 50 || res.V.Issues[1].Number != 48 { + t.Errorf("kept %d,%d, want 50,48 (the row behind the malformed one still collected)", res.V.Issues[0].Number, res.V.Issues[1].Number) } } @@ -384,8 +385,8 @@ func TestListUpdatedIssuesZeroSinceWalksAll(t *testing.T) { } // A zero since is never strictly greater than any updated_at, so nothing // stops the walk short — both pages collected, walk ends at no rel="next". - if len(res.V) != 2 { - t.Fatalf("kept %d (%+v), want 2 (zero since walks to the last page)", len(res.V), res.V) + if len(res.V.Issues) != 2 { + t.Fatalf("kept %d (%+v), want 2 (zero since walks to the last page)", len(res.V.Issues), res.V.Issues) } if rt.calls != 2 { t.Errorf("calls = %d, want 2 (walk ended only at no rel=\"next\")", rt.calls) @@ -395,7 +396,7 @@ func TestListUpdatedIssuesZeroSinceWalksAll(t *testing.T) { // TestListUpdatedIssuesFiltersPRs: /repos/{repo}/issues interleaves PR rows // (pull_request marker); the updated-order walk drops them, keeping only // issue-shaped rows, on the updated-order endpoint. -func TestListUpdatedIssuesFiltersPRs(t *testing.T) { +func TestListUpdatedIssuesSplitsPRs(t *testing.T) { body := `[ {"number":45,"state":"open","html_url":"u45","updated_at":"2026-08-01T12:00:02Z","pull_request":{"url":"pr"}}, {"number":44,"state":"open","html_url":"u44","updated_at":"2026-08-01T12:00:01Z"} @@ -409,8 +410,12 @@ func TestListUpdatedIssuesFiltersPRs(t *testing.T) { if err != nil { t.Fatalf("ListUpdatedIssues: %v", err) } - if len(res.V) != 1 || res.V[0].Number != 44 { - t.Fatalf("kept %+v, want only issue #44 (PR #45 filtered)", res.V) + if len(res.V.Issues) != 1 || res.V.Issues[0].Number != 44 { + t.Fatalf("issues %+v, want only issue #44", res.V.Issues) + } + want := []UpdatedPull{{Number: 45, State: "open", UpdatedAt: time.Date(2026, 8, 1, 12, 0, 2, 0, time.UTC)}} + if !slices.Equal(res.V.Pulls, want) { + t.Fatalf("pulls %+v, want %+v (PR rows kept in order beside the issues)", res.V.Pulls, want) } if got := rt.requests[0].URL.Path; got != "/repos/org/repo/issues" { t.Errorf("path = %q, want the issues endpoint", got) diff --git a/go/internal/forge/provider.go b/go/internal/forge/provider.go index 3db886a10..c5272f964 100644 --- a/go/internal/forge/provider.go +++ b/go/internal/forge/provider.go @@ -110,6 +110,17 @@ type PullRequest struct { Reviews []Review // Threads are the inline (and PR-level) review comment threads. Threads []ReviewThread + // CreatedAt and UpdatedAt are the forge's timestamps; zero when unparseable. + CreatedAt time.Time + UpdatedAt time.Time + // ClosingRefs are the issues the forge says this PR closes; only a full read fills them. + ClosingRefs []IssueRef +} + +// IssueRef names an issue by repo ("owner/name") and number on the PR's forge. +type IssueRef struct { + Repo string + Number uint64 } // ChangedStats is the diff-size roll-up of a pull request. diff --git a/go/internal/forge/pull_link_reads_test.go b/go/internal/forge/pull_link_reads_test.go new file mode 100644 index 000000000..03379c618 --- /dev/null +++ b/go/internal/forge/pull_link_reads_test.go @@ -0,0 +1,123 @@ +package forge + +// Forge reads the PR-to-issue linkage depends on: PR timestamps, closing +// references on the first GraphQL page only, and the open-PR list. + +import ( + "context" + "slices" + "strings" + "testing" + "time" +) + +var ( + prCreated = time.Date(2026, 9, 1, 10, 0, 0, 0, time.UTC) + prUpdated = time.Date(2026, 9, 2, 11, 30, 0, 0, time.UTC) +) + +func TestCreatePullRequestDecodesTimestamps(t *testing.T) { + rt := &scriptedRoundTripper{responses: []scriptedResponse{{status: 201, body: `{"number":13,"state":"open", + "created_at":"2026-09-01T10:00:00Z","updated_at":"2026-09-02T11:30:00Z","head":{"ref":"f"},"base":{"ref":"main"}}`}}} + g := newTestGitHub(rt, &fakeTokenSource{token: "t"}) + + got, err := g.CreatePullRequest(context.Background(), "org/repo", CreatePR{Title: "x", Body: "b", HeadRef: "f"}) + if err != nil { + t.Fatalf("CreatePullRequest: %v", err) + } + if !got.CreatedAt.Equal(prCreated) || !got.UpdatedAt.Equal(prUpdated) { + t.Fatalf("times = %v / %v, want %v / %v", got.CreatedAt, got.UpdatedAt, prCreated, prUpdated) + } +} + +// closingRefs is one page of closingIssuesReferences. +func closingRefs(refs ...IssueRef) *ghGQLClosingRefs { + c := &ghGQLClosingRefs{} + for _, r := range refs { + n := ghGQLClosingRef{Number: r.Number} + n.Repository.NameWithOwner = r.Repo + c.Nodes = append(c.Nodes, n) + } + return c +} + +// GetPullRequest decodes the detail timestamps and reads closing references on +// the first GraphQL page only, not again on a thread page. +func TestGetPullRequestClosingRefsFirstPageOnly(t *testing.T) { + refs := []IssueRef{{Repo: "org/repo", Number: 3}, {Repo: "Other/Repo", Number: 9}} + page1 := gqlBody(t, ghGQLPullData{Repository: &ghGQLRepo{PullRequest: &ghGQLPull{ + ReviewThreads: threadsConn(true, "CUR1", thread("T1", "a.go", false, comment("x", "User", "one"))), + Commits: rollup(headSHA, false, "c"), + ClosingRefs: closingRefs(refs...), + }}}) + page2 := pullPage(t, threadsConn(false, "CUR2"), nil) + rt := &scriptedRoundTripper{responses: []scriptedResponse{ + ok200(`{"number":7,"state":"open","created_at":"2026-09-01T10:00:00Z","updated_at":"2026-09-02T11:30:00Z", + "head":{"ref":"f","sha":"` + headSHA + `"},"base":{"ref":"main"},"user":{"login":"a"}}`), + ok200(`[]`), ok200(`{"check_runs": []}`), ok200(`{"statuses": []}`), + ok200(page1), ok200(page2), + }} + g := newTestGitHub(rt, &fakeTokenSource{token: "t"}) + + got, err := g.GetPullRequest(context.Background(), "org/repo", 7) + if err != nil { + t.Fatalf("GetPullRequest: %v", err) + } + if !slices.Equal(got.ClosingRefs, refs) { + t.Fatalf("ClosingRefs = %+v, want %+v", got.ClosingRefs, refs) + } + if !got.CreatedAt.Equal(prCreated) || !got.UpdatedAt.Equal(prUpdated) { + t.Fatalf("times = %v / %v", got.CreatedAt, got.UpdatedAt) + } + if b := readReqBody(t, rt.requests[4]); !strings.Contains(b, `"refs":true`) { + t.Errorf("page 1 must request closing refs: %s", b) + } + if b := readReqBody(t, rt.requests[5]); !strings.Contains(b, `"refs":false`) { + t.Errorf("page 2 must not re-request closing refs: %s", b) + } +} + +// The checks-only read never asks for closing references. +func TestChecksSkipsClosingRefs(t *testing.T) { + rt := &scriptedRoundTripper{responses: pullRESTLegs("", ok200(noRollupGraphQL(headSHA)))} + rt.responses = append(rt.responses[:1], rt.responses[2:]...) // Checks reads no reviews + g := newTestGitHub(rt, &fakeTokenSource{token: "t"}) + + if _, err := g.Checks(context.Background(), "org/repo", 7); err != nil { + t.Fatalf("Checks: %v", err) + } + if b := readReqBody(t, rt.requests[len(rt.requests)-1]); !strings.Contains(b, `"refs":false`) { + t.Errorf("checks read requested closing refs: %s", b) + } +} + +// A GraphQL error on the closing-refs page fails the read rather than dropping links. +func TestGetPullRequestGraphQLErrorIsError(t *testing.T) { + rt := &scriptedRoundTripper{responses: pullRESTLegs("", ok200(`{"errors":[{"message":"boom"}]}`))} + g := newTestGitHub(rt, &fakeTokenSource{token: "t"}) + + if _, err := g.GetPullRequest(context.Background(), "org/repo", 7); err == nil { + t.Fatal("GetPullRequest succeeded on a GraphQL error") + } +} + +func TestListOpenPullRequests(t *testing.T) { + rt := &scriptedRoundTripper{responses: []scriptedResponse{ + {status: 200, body: `[{"number":9,"state":"open","updated_at":"2026-09-02T11:30:00Z"}]`, + headers: map[string]string{"Link": `; rel="next"`}}, + ok200(`[{"number":4,"state":"open","updated_at":"2026-09-01T10:00:00Z"}]`), + }} + g := newTestGitHub(rt, &fakeTokenSource{token: "t"}) + + got, err := g.ListOpenPullRequests(context.Background(), "org/repo") + if err != nil { + t.Fatalf("ListOpenPullRequests: %v", err) + } + want := []UpdatedPull{{Number: 9, State: "open", UpdatedAt: prUpdated}, {Number: 4, State: "open", UpdatedAt: prCreated}} + if !slices.Equal(got, want) { + t.Fatalf("got %+v, want %+v", got, want) + } + if q := rt.requests[0].URL.Query(); rt.requests[0].URL.Path != "/repos/org/repo/pulls" || q.Get("state") != "open" { + t.Errorf("request = %s, want the open pulls list", rt.requests[0].URL) + } +} diff --git a/go/internal/forge/testdata/github/create_pull_request.json b/go/internal/forge/testdata/github/create_pull_request.json index 7ff2f54e6..f345f3162 100644 --- a/go/internal/forge/testdata/github/create_pull_request.json +++ b/go/internal/forge/testdata/github/create_pull_request.json @@ -46,7 +46,10 @@ "Changed": { "Files": 0, "Additions": 0, "Deletions": 0 }, "Checks": { "HeadSHA": "", "State": "", "Checks": null }, "Reviews": null, - "Threads": null + "Threads": null, + "CreatedAt": "0001-01-01T00:00:00Z", + "UpdatedAt": "0001-01-01T00:00:00Z", + "ClosingRefs": null } } } diff --git a/go/internal/forge/testdata/github/get_pull_request.json b/go/internal/forge/testdata/github/get_pull_request.json index a5c8bd8bb..d40b3dbe9 100644 --- a/go/internal/forge/testdata/github/get_pull_request.json +++ b/go/internal/forge/testdata/github/get_pull_request.json @@ -71,6 +71,9 @@ "data": { "repository": { "pullRequest": { + "closingIssuesReferences": { + "nodes": [] + }, "reviewThreads": { "pageInfo": { "hasNextPage": false, @@ -218,7 +221,10 @@ } ] } - ] + ], + "CreatedAt": "0001-01-01T00:00:00Z", + "UpdatedAt": "0001-01-01T00:00:00Z", + "ClosingRefs": null } } } diff --git a/go/internal/forge/testdata/github/transition_pull_request_close.json b/go/internal/forge/testdata/github/transition_pull_request_close.json index f3625d67a..4d28a27f0 100644 --- a/go/internal/forge/testdata/github/transition_pull_request_close.json +++ b/go/internal/forge/testdata/github/transition_pull_request_close.json @@ -47,7 +47,10 @@ "Changed": { "Files": 3, "Additions": 10, "Deletions": 2 }, "Checks": { "HeadSHA": "", "State": "", "Checks": null }, "Reviews": null, - "Threads": null + "Threads": null, + "CreatedAt": "0001-01-01T00:00:00Z", + "UpdatedAt": "0001-01-01T00:00:00Z", + "ClosingRefs": null } } } diff --git a/go/internal/forge/testdata/github/transition_pull_request_reopen.json b/go/internal/forge/testdata/github/transition_pull_request_reopen.json index f7626a2cd..20afae2cb 100644 --- a/go/internal/forge/testdata/github/transition_pull_request_reopen.json +++ b/go/internal/forge/testdata/github/transition_pull_request_reopen.json @@ -47,7 +47,10 @@ "Changed": { "Files": 3, "Additions": 10, "Deletions": 2 }, "Checks": { "HeadSHA": "", "State": "", "Checks": null }, "Reviews": null, - "Threads": null + "Threads": null, + "CreatedAt": "0001-01-01T00:00:00Z", + "UpdatedAt": "0001-01-01T00:00:00Z", + "ClosingRefs": null } } } diff --git a/go/internal/ingest/board_reconcile.go b/go/internal/ingest/board_reconcile.go index 18631e76a..66b212215 100644 --- a/go/internal/ingest/board_reconcile.go +++ b/go/internal/ingest/board_reconcile.go @@ -43,7 +43,7 @@ type BoardStore interface { // forge.GitHub.ListUpdatedIssues at T5. LOCAL + structural: this package never // imports the concrete provider. type updatedLister interface { - ListUpdatedIssues(ctx context.Context, repo string, since time.Time, etag string) (forge.ConditionalResult[[]forge.Issue], error) + ListUpdatedIssues(ctx context.Context, repo string, since time.Time, etag string) (forge.ConditionalResult[forge.UpdatedRows], error) } // BoardReconcileConfig configures the board reconciliation sweep. @@ -181,7 +181,7 @@ func (rc *BoardReconciler) reconcileRepo(ctx context.Context, repo string) error var maxMark time.Time poison := 0 - for _, row := range res.V { + for _, row := range res.V.Issues { if ctx.Err() != nil { return ctx.Err() } diff --git a/go/internal/ingest/board_reconcile_test.go b/go/internal/ingest/board_reconcile_test.go index de38ff643..be9f6500f 100644 --- a/go/internal/ingest/board_reconcile_test.go +++ b/go/internal/ingest/board_reconcile_test.go @@ -29,19 +29,24 @@ type fakeUpdatedLister struct { calls atomic.Int64 } -func (l *fakeUpdatedLister) ListUpdatedIssues(_ context.Context, _ string, _ time.Time, _ string) (forge.ConditionalResult[[]forge.Issue], error) { +func (l *fakeUpdatedLister) ListUpdatedIssues(_ context.Context, _ string, _ time.Time, _ string) (forge.ConditionalResult[forge.UpdatedRows], error) { n := l.calls.Add(1) if l.err != nil { - return forge.ConditionalResult[[]forge.Issue]{}, l.err + return forge.ConditionalResult[forge.UpdatedRows]{}, l.err } if len(l.results) == 0 { - return forge.ConditionalResult[[]forge.Issue]{NotModified: true}, nil + return forge.ConditionalResult[forge.UpdatedRows]{NotModified: true}, nil } i := int(n) - 1 if i >= len(l.results) { i = len(l.results) - 1 } - return l.results[i], nil + return asUpdatedRows(l.results[i]), nil +} + +// asUpdatedRows wraps a scripted issue-only result as an updated-order walk. +func asUpdatedRows(r forge.ConditionalResult[[]forge.Issue]) forge.ConditionalResult[forge.UpdatedRows] { + return forge.ConditionalResult[forge.UpdatedRows]{V: forge.UpdatedRows{Issues: r.V}, ETag: r.ETag, NotModified: r.NotModified} } // storedMark is one repo's persisted watermark row. @@ -250,7 +255,7 @@ type sinceAwareLister struct { lastWindow atomic.Int64 } -func (l *sinceAwareLister) ListUpdatedIssues(_ context.Context, _ string, since time.Time, _ string) (forge.ConditionalResult[[]forge.Issue], error) { +func (l *sinceAwareLister) ListUpdatedIssues(_ context.Context, _ string, since time.Time, _ string) (forge.ConditionalResult[forge.UpdatedRows], error) { l.calls.Add(1) var out []forge.Issue for _, iss := range l.all { @@ -260,7 +265,7 @@ func (l *sinceAwareLister) ListUpdatedIssues(_ context.Context, _ string, since out = append(out, iss) } l.lastWindow.Store(int64(len(out))) - return forge.ConditionalResult[[]forge.Issue]{V: out, ETag: l.etag}, nil + return forge.ConditionalResult[forge.UpdatedRows]{V: forge.UpdatedRows{Issues: out}, ETag: l.etag}, nil } // TestBoardSweepPoisonNewestBoundedAcrossSweeps: a persistently-failing row that @@ -375,14 +380,14 @@ type perRepoLister struct { results map[string]forge.ConditionalResult[[]forge.Issue] } -func (l *perRepoLister) ListUpdatedIssues(_ context.Context, repo string, _ time.Time, _ string) (forge.ConditionalResult[[]forge.Issue], error) { +func (l *perRepoLister) ListUpdatedIssues(_ context.Context, repo string, _ time.Time, _ string) (forge.ConditionalResult[forge.UpdatedRows], error) { if err, ok := l.errRepos[repo]; ok { - return forge.ConditionalResult[[]forge.Issue]{}, err + return forge.ConditionalResult[forge.UpdatedRows]{}, err } if res, ok := l.results[repo]; ok { - return res, nil + return asUpdatedRows(res), nil } - return forge.ConditionalResult[[]forge.Issue]{NotModified: true}, nil + return forge.ConditionalResult[forge.UpdatedRows]{NotModified: true}, nil } // TestBoardRunStartupSweepFiresImmediately: Run performs one immediate sweep at diff --git a/go/internal/store/db/forge_cursors.sql.go b/go/internal/store/db/forge_cursors.sql.go index 1ed9a6324..f81be7f3d 100644 --- a/go/internal/store/db/forge_cursors.sql.go +++ b/go/internal/store/db/forge_cursors.sql.go @@ -150,6 +150,49 @@ func (q *Queries) LoadForgeRepoWatermark(ctx context.Context, arg LoadForgeRepoW return i, err } +const loadPRsBackfilledAt = `-- name: LoadPRsBackfilledAt :one +SELECT prs_backfilled_at FROM forge_repo_subscriptions + WHERE forge_provider = $1 AND forge_host = $2 AND repo = $3 +` + +type LoadPRsBackfilledAtParams struct { + ForgeProvider int16 + ForgeHost string + Repo string +} + +func (q *Queries) LoadPRsBackfilledAt(ctx context.Context, arg LoadPRsBackfilledAtParams) (pgtype.Timestamptz, error) { + row := q.db.QueryRow(ctx, loadPRsBackfilledAt, arg.ForgeProvider, arg.ForgeHost, arg.Repo) + var prs_backfilled_at pgtype.Timestamptz + err := row.Scan(&prs_backfilled_at) + return prs_backfilled_at, err +} + +const markPRsBackfilled = `-- name: MarkPRsBackfilled :execrows +UPDATE forge_repo_subscriptions SET prs_backfilled_at = $4 + WHERE forge_provider = $1 AND forge_host = $2 AND repo = $3 +` + +type MarkPRsBackfilledParams struct { + ForgeProvider int16 + ForgeHost string + Repo string + PrsBackfilledAt pgtype.Timestamptz +} + +func (q *Queries) MarkPRsBackfilled(ctx context.Context, arg MarkPRsBackfilledParams) (int64, error) { + result, err := q.db.Exec(ctx, markPRsBackfilled, + arg.ForgeProvider, + arg.ForgeHost, + arg.Repo, + arg.PrsBackfilledAt, + ) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} + const setForgeRepoSubscriptionEnabled = `-- name: SetForgeRepoSubscriptionEnabled :execrows UPDATE forge_repo_subscriptions SET enabled = $4 diff --git a/go/internal/store/db/pull_requests.sql.go b/go/internal/store/db/pull_requests.sql.go index 2149fb749..6f41abb2d 100644 --- a/go/internal/store/db/pull_requests.sql.go +++ b/go/internal/store/db/pull_requests.sql.go @@ -221,6 +221,30 @@ func (q *Queries) ListPullRequestLinks(ctx context.Context, arg ListPullRequestL return items, nil } +const pullRequestForgeUpdatedAt = `-- name: PullRequestForgeUpdatedAt :one +SELECT forge_updated_at FROM pull_requests + WHERE forge_provider = $1 AND forge_host = $2 AND repo = $3 AND number = $4 +` + +type PullRequestForgeUpdatedAtParams struct { + ForgeProvider int16 + ForgeHost string + Repo string + Number int64 +} + +func (q *Queries) PullRequestForgeUpdatedAt(ctx context.Context, arg PullRequestForgeUpdatedAtParams) (pgtype.Timestamptz, error) { + row := q.db.QueryRow(ctx, pullRequestForgeUpdatedAt, + arg.ForgeProvider, + arg.ForgeHost, + arg.Repo, + arg.Number, + ) + var forge_updated_at pgtype.Timestamptz + err := row.Scan(&forge_updated_at) + return forge_updated_at, err +} + const pullRequestsForIssues = `-- name: PullRequestsForIssues :many WITH want AS ( SELECT unnest($1::smallint[]) AS provider, diff --git a/go/internal/store/db/querier.go b/go/internal/store/db/querier.go index 5ae39b1a7..00260ace8 100644 --- a/go/internal/store/db/querier.go +++ b/go/internal/store/db/querier.go @@ -393,6 +393,7 @@ type Querier interface { // read methods map the generated rows (provider int, nullable swept_updated_at) // back to the domain time.Time / ForgeRepoSubscription. LoadForgeRepoWatermark(ctx context.Context, arg LoadForgeRepoWatermarkParams) (LoadForgeRepoWatermarkRow, error) + LoadPRsBackfilledAt(ctx context.Context, arg LoadPRsBackfilledAtParams) (pgtype.Timestamptz, error) // Channel-pins (pinned board) queries (sqlc adoption T3, RIG-3034). These // replace the inline SQL literals in internal/store/channel_pins.go; the // hand-written Store methods and the in-tx FOR UPDATE lock / cap-check control @@ -465,6 +466,7 @@ type Querier interface { // an append would count the append's events twice or lose them. LockTokenUsage(ctx context.Context) error MarkMentionsRouted(ctx context.Context, arg MarkMentionsRoutedParams) error + MarkPRsBackfilled(ctx context.Context, arg MarkPRsBackfilledParams) (int64, error) MergeTopicLastSeq(ctx context.Context, arg MergeTopicLastSeqParams) error MessageByID(ctx context.Context, id string) (MessageByIDRow, error) MessageChannel(ctx context.Context, id string) (string, error) @@ -479,6 +481,7 @@ type Querier interface { PinnedEntries(ctx context.Context, channelID string) ([]PinnedEntriesRow, error) PlacementForAgent(ctx context.Context, agentAccountID string) (PlacementForAgentRow, error) PruneTranscriptEntries(ctx context.Context, arg PruneTranscriptEntriesParams) error + PullRequestForgeUpdatedAt(ctx context.Context, arg PullRequestForgeUpdatedAtParams) (pgtype.Timestamptz, error) // Explicit links attach directly; a closing reference attaches only when none // of the PR's explicit targets is an issue on the board. PullRequestsForIssues(ctx context.Context, arg PullRequestsForIssuesParams) ([]PullRequestsForIssuesRow, error) diff --git a/go/internal/store/forge_cursors.go b/go/internal/store/forge_cursors.go index 21244d147..b353d33a6 100644 --- a/go/internal/store/forge_cursors.go +++ b/go/internal/store/forge_cursors.go @@ -202,3 +202,47 @@ func (s *Store) SetForgeRepoSubscriptionEnabled(ctx context.Context, provider Fo } return nil } + +// PRsBackfilledAt reads when the repo's bounded PR backfill last ran. ok is +// false when it never ran or the repo has no subscription row. +func (s *Store) PRsBackfilledAt(ctx context.Context, provider ForgeProvider, host, repo string) (time.Time, bool, error) { + if err := validCoordinate(provider, host, repo); err != nil { + return time.Time{}, false, err + } + at, err := s.q.LoadPRsBackfilledAt(ctx, db.LoadPRsBackfilledAtParams{ + ForgeProvider: int16(provider), //nolint:gosec // G115: ForgeProvider is a CHECK-constrained 1..4 enum, always within int16 + ForgeHost: host, + Repo: repo, + }) + if noRows(err) { + return time.Time{}, false, nil + } + if err != nil { + return time.Time{}, false, fmt.Errorf("store: load prs backfilled at: %w", err) + } + return at.Time, at.Valid, nil +} + +// MarkPRsBackfilled records that the repo's PR backfill ran at at. An unknown +// repo is ErrNotFound. +func (s *Store) MarkPRsBackfilled(ctx context.Context, provider ForgeProvider, host, repo string, at time.Time) error { + if err := validCoordinate(provider, host, repo); err != nil { + return err + } + if at.IsZero() { + return fmt.Errorf("%w: backfill time is required", ErrInvalidArgument) + } + affected, err := s.q.MarkPRsBackfilled(ctx, db.MarkPRsBackfilledParams{ + ForgeProvider: int16(provider), //nolint:gosec // G115: ForgeProvider is a CHECK-constrained 1..4 enum, always within int16 + ForgeHost: host, + Repo: repo, + PrsBackfilledAt: pgtype.Timestamptz{Time: at, Valid: true}, + }) + if err != nil { + return fmt.Errorf("store: mark prs backfilled: %w", err) + } + if affected == 0 { + return fmt.Errorf("%w: forge repo subscription (%d, %q, %q)", ErrNotFound, provider, host, repo) + } + return nil +} diff --git a/go/internal/store/pull_requests.go b/go/internal/store/pull_requests.go index 7cf9421dd..cbbdf0f6f 100644 --- a/go/internal/store/pull_requests.go +++ b/go/internal/store/pull_requests.go @@ -322,3 +322,23 @@ func (s *Store) FallbackIssuesForTarget(ctx context.Context, issue ForgeCoord) ( } return out, nil } + +// PullRequestUpdatedAt reads the stored forge_updated_at of pr; ok is false when +// the PR was never stored. The sweep hydrates only rows newer than this. +func (s *Store) PullRequestUpdatedAt(ctx context.Context, pr ForgeCoord) (time.Time, bool, error) { + pr = pr.normalized() + if err := pr.valid(); err != nil { + return time.Time{}, false, err + } + p := pr.db() + at, err := s.q.PullRequestForgeUpdatedAt(ctx, db.PullRequestForgeUpdatedAtParams{ + ForgeProvider: p.provider, ForgeHost: p.host, Repo: p.repo, Number: p.number, + }) + if noRows(err) { + return time.Time{}, false, nil + } + if err != nil { + return time.Time{}, false, fmt.Errorf("store: pull request updated at: %w", err) + } + return at.Time, true, nil +} diff --git a/go/internal/store/pull_requests_pgtest_test.go b/go/internal/store/pull_requests_pgtest_test.go index 8188032a1..5e5cb7da7 100644 --- a/go/internal/store/pull_requests_pgtest_test.go +++ b/go/internal/store/pull_requests_pgtest_test.go @@ -301,3 +301,37 @@ func sameCoords(got, want []ForgeCoord) bool { return !slices.Contains(got, c) }) } + +func TestPullRequestUpdatedAtGate(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + if _, ok, err := s.PullRequestUpdatedAt(ctx, ghCoord("a/b", 10)); err != nil || ok { + t.Fatalf("unknown PR = ok %v, err %v; want not found", ok, err) + } + mustUpsertPR(t, ctx, s, prRow("A/b", 10, "open", prBase, prBase.Add(time.Hour))) + at, ok, err := s.PullRequestUpdatedAt(ctx, ghCoord("a/B", 10)) + if err != nil || !ok || !at.Equal(prBase.Add(time.Hour)) { + t.Fatalf("PullRequestUpdatedAt = %v, %v, %v; want the stored forge time", at, ok, err) + } +} + +func TestPRsBackfilledMark(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + if err := s.EnsureForgeRepoSubscription(ctx, ForgeRepoSubscription{Provider: ForgeProviderGitHub, Host: "github.com", Repo: "a/b"}); err != nil { + t.Fatalf("EnsureForgeRepoSubscription: %v", err) + } + if _, ok, err := s.PRsBackfilledAt(ctx, ForgeProviderGitHub, "github.com", "a/b"); err != nil || ok { + t.Fatalf("fresh repo backfilled = %v, %v; want NULL", ok, err) + } + if err := s.MarkPRsBackfilled(ctx, ForgeProviderGitHub, "github.com", "a/b", prBase); err != nil { + t.Fatalf("MarkPRsBackfilled: %v", err) + } + at, ok, err := s.PRsBackfilledAt(ctx, ForgeProviderGitHub, "github.com", "a/b") + if err != nil || !ok || !at.Equal(prBase) { + t.Fatalf("PRsBackfilledAt = %v, %v, %v; want %v", at, ok, err, prBase) + } + if err := s.MarkPRsBackfilled(ctx, ForgeProviderGitHub, "github.com", "x/y", prBase); !errors.Is(err, ErrNotFound) { + t.Fatalf("unknown repo err = %v, want ErrNotFound", err) + } +} diff --git a/go/internal/store/queries/forge_cursors.sql b/go/internal/store/queries/forge_cursors.sql index c187cc969..2a10ffb88 100644 --- a/go/internal/store/queries/forge_cursors.sql +++ b/go/internal/store/queries/forge_cursors.sql @@ -42,3 +42,11 @@ ORDER BY repo ASC; UPDATE forge_repo_subscriptions SET enabled = $4 WHERE forge_provider = $1 AND forge_host = $2 AND repo = $3; + +-- name: LoadPRsBackfilledAt :one +SELECT prs_backfilled_at FROM forge_repo_subscriptions + WHERE forge_provider = $1 AND forge_host = $2 AND repo = $3; + +-- name: MarkPRsBackfilled :execrows +UPDATE forge_repo_subscriptions SET prs_backfilled_at = $4 + WHERE forge_provider = $1 AND forge_host = $2 AND repo = $3; diff --git a/go/internal/store/queries/pull_requests.sql b/go/internal/store/queries/pull_requests.sql index 3dc3cf8d0..5b57305ab 100644 --- a/go/internal/store/queries/pull_requests.sql +++ b/go/internal/store/queries/pull_requests.sql @@ -91,3 +91,7 @@ SELECT DISTINCT c.issue_forge_provider, c.issue_forge_host, c.issue_repo, c.issu AND e.issue_forge_provider = $1 AND e.issue_forge_host = $2 AND e.issue_repo = $3 AND e.issue_number = $4 ORDER BY c.issue_forge_provider, c.issue_forge_host, c.issue_repo, c.issue_number; + +-- name: PullRequestForgeUpdatedAt :one +SELECT forge_updated_at FROM pull_requests + WHERE forge_provider = $1 AND forge_host = $2 AND repo = $3 AND number = $4;