diff --git a/go/internal/board/issue_projection.go b/go/internal/board/issue_projection.go index 470b08666..218b01a54 100644 --- a/go/internal/board/issue_projection.go +++ b/go/internal/board/issue_projection.go @@ -4,6 +4,7 @@ package board import ( "context" + "errors" "fmt" "sort" "sync" @@ -28,8 +29,13 @@ type IssueProjection struct { bus *events.Bus[busPayload] store *store.Store - mu sync.RWMutex - issues map[string]*compassv1.Issue // id -> latest canonical issue (in-memory cache) + // writeMu orders each store read with the cache write it feeds, so a slower + // writer never caches an older row or prs list over a newer one. Never held + // under mu; taken before it. + writeMu sync.Mutex + mu sync.RWMutex + issues map[string]*compassv1.Issue // id -> latest canonical issue (in-memory cache) + byCoord map[store.ForgeCoord]string // normalized forge coordinate -> issue id } // NewIssueProjection constructs an empty board over the SubscribeEvents bus it @@ -37,9 +43,10 @@ type IssueProjection struct { // through. Rehydrate seeds the map from the store before serving. func NewIssueProjection(bus *events.Bus[busPayload], st *store.Store) *IssueProjection { return &IssueProjection{ - bus: bus, - store: st, - issues: make(map[string]*compassv1.Issue), + bus: bus, + store: st, + issues: make(map[string]*compassv1.Issue), + byCoord: make(map[store.ForgeCoord]string), } } @@ -48,23 +55,14 @@ func NewIssueProjection(bus *events.Bus[busPayload], st *store.Store) *IssueProj // commit; returns the stable id), (3) GetIssue(id) to read back the FULL row // (forge fields just written + the store-owned state/machinery, so the cached + // fanned Issue reflects committed truth incl. a prior human-set state), (4) map -// store.Issue -> *compassv1.Issue, (5) record in the map + Publish the issue=16 -// variant, atomic under mu. Returns error on any store failure (part 3 stops on it). +// store.Issue -> *compassv1.Issue with its prs, (5) record in the map + Publish +// the issue=16 variant, atomic under mu, then (6) republish the closing-ref +// issues of PRs that now attach here. Returns error on any store failure. // -// Lock discipline (load-bearing, mirrors PublishSessionStatus): the DURABLE PG -// commit (UpsertIssueForgeFields + the GetIssue read-back) happens BEFORE taking -// p.mu — the DB round-trip must NOT be held under the projection mutex, which -// would serialize all ingestion on the database. Only the map-record + the -// bus.Publish run under the lock, where Publish is non-blocking (a per-subscriber -// select/default under the bus's own distinct mutex), so record and fan-out are -// atomic to any Snapshot reader and deadlock-free. -// -// This assumes single-threaded ingestion per coordinate (the part-3 poller): -// the lock is released between the DB commit and the map-record, so two -// concurrent publishes of the SAME coordinate could record in commit order or -// in lock order. Harmless while one poller owns each coordinate; a -// per-coordinate guard would be needed if concurrent same-coordinate ingestion -// is ever introduced. +// Lock discipline: the upsert runs unlocked. writeMu, taken before mu, spans the +// read-back through the record, so every projection writer is serialized and a +// write committed meanwhile is either read here or cached after this. mu covers +// only the map-record and the non-blocking bus.Publish, never a DB round-trip. func (p *IssueProjection) PublishIssueUpdate(ctx context.Context, issue *compassv1.Issue) error { // (1)-(2) durable commit at the forge coordinate; the returned id is stable // across re-polls (the coordinate is the idempotency key). @@ -72,6 +70,8 @@ func (p *IssueProjection) PublishIssueUpdate(ctx context.Context, issue *compass if err != nil { return fmt.Errorf("board: upsert issue forge fields: %w", err) } + p.writeMu.Lock() + defer p.writeMu.Unlock() // (3) read back the FULL committed row: the forge fields just written PLUS // the store-owned state/machinery, so the fanned Issue reflects committed // truth — including a prior human-set lifecycle state the forge re-poll did @@ -80,45 +80,75 @@ func (p *IssueProjection) PublishIssueUpdate(ctx context.Context, issue *compass if err != nil { return fmt.Errorf("board: read back committed issue: %w", err) } - // (4) map committed store.Issue -> wire Issue OUTSIDE the lock. - wire := issueToProto(committed) + // (4) map committed store.Issue -> wire Issue with its prs, OUTSIDE the lock. + wires, err := p.IssueToProtoWithPrs(ctx, committed) + if err != nil { + return err + } + wire := wires[0] + coord := storeIssueCoord(committed) // (5) record + fan out atomically under the write lock. p.mu.Lock() - defer p.mu.Unlock() - p.issues[wire.GetId()] = wire + p.record(wire, coord) p.bus.Publish(&compassv1.SubscribeEventsResponse{ Payload: &compassv1.SubscribeEventsResponse_Issue{Issue: wire}, }) - return nil + p.mu.Unlock() + + // PRs that fell back to closing refs while this issue was off the board now + // attach here; republish those closing-ref issues without them. + fallback, err := p.store.FallbackIssuesForTarget(ctx, coord) + if err != nil { + return fmt.Errorf("board: fallback issues for target: %w", err) + } + return p.republishPrsLocked(ctx, fallback) } -// RecordAndPublish is the STATE-ONLY record+publish the write-path transition -// executor (server/board.go, agent primary lifecycle T3-a) drives after it has -// already committed the new canonical state to Postgres AND read the full row -// back. It is the step-(5) tail of PublishIssueUpdate WITHOUT the store -// upsert/read-back: the executor owns the durable commit (a forge-only upsert -// would demand forge fields and could not carry the state column), so this only -// maps the committed row to the wire Issue and records+fans it — never a store -// write. committed is the executor's read-back of committed truth. -// -// Lock discipline mirrors PublishIssueUpdate exactly: the map -> proto mapping -// runs OUTSIDE p.mu (no DB round-trip is held under the projection mutex; here -// there is no DB work at all), and only the map-record + the non-blocking -// bus.Publish run under the lock, so record and fan-out are atomic to any -// Snapshot reader and deadlock-free. Per-coordinate serialization is the -// executor's job (its per-issue transition lock), not this projection's. -func (p *IssueProjection) RecordAndPublish(committed store.Issue) { - // Map committed store.Issue -> wire Issue OUTSIDE the lock. +// RecordAndPublish records and fans out an issue the transition executor +// (server/board.go) has already committed; it never writes the store. Under +// writeMu it re-reads the row, since a forge upsert may have landed after the +// caller's read, and loads prs for an uncached issue. Those reads are +// best-effort: on failure it still publishes committed, with cached or no prs, +// and returns the error for the caller to log, so a durable transition always +// reaches the board. +func (p *IssueProjection) RecordAndPublish(ctx context.Context, committed store.Issue) error { + p.writeMu.Lock() + defer p.writeMu.Unlock() + var readErr error + if p.store != nil { + fresh, err := p.store.GetIssue(ctx, committed.ID) + if err == nil { + committed = fresh + } else { + readErr = fmt.Errorf("board: read back transitioned issue: %w", err) + } + } wire := issueToProto(committed) - // Record + fan out atomically under the write lock. + // A state change never touches links, so cached prs carry over. + p.mu.RLock() + prev, cached := p.issues[wire.GetId()] + p.mu.RUnlock() + if cached { + wire.Prs = prev.GetPrs() + } else if p.store != nil { + coord := storeIssueCoord(committed) + prs, err := p.loadPrs(ctx, []store.ForgeCoord{coord}) + if err == nil { + wire.Prs = prs[coord] + } else { + readErr = errors.Join(readErr, err) + } + } + p.mu.Lock() defer p.mu.Unlock() - p.issues[wire.GetId()] = wire + p.record(wire, storeIssueCoord(committed)) p.bus.Publish(&compassv1.SubscribeEventsResponse{ Payload: &compassv1.SubscribeEventsResponse_Issue{Issue: wire}, }) + return readErr } // Snapshot returns every issue on the board (all states incl. ARCHIVED — the @@ -142,34 +172,35 @@ func (p *IssueProjection) Snapshot() []*compassv1.Issue { // subscribed yet at boot); it seeds the map so the first Snapshot/fan-out is // complete. func (p *IssueProjection) Rehydrate(ctx context.Context) error { + p.writeMu.Lock() + defer p.writeMu.Unlock() rows, err := p.store.ListIssues(ctx) if err != nil { return fmt.Errorf("board: rehydrate issues: %w", err) } + wires, err := p.IssueToProtoWithPrs(ctx, rows...) + if err != nil { + return fmt.Errorf("board: rehydrate issues: %w", err) + } p.mu.Lock() defer p.mu.Unlock() - for _, si := range rows { - proto := issueToProto(si) - p.issues[proto.GetId()] = proto + for i, si := range rows { + p.record(wires[i], storeIssueCoord(si)) } return nil } -// IssueToProto maps a store-native issue to the canonical wire Issue. It is the -// exported form of issueToProto for the write-path transition executor -// (server/board.go), which reads a committed row back through the store and must -// return it on the SetIssueState response wire — the store<->wire mapping stays -// owned by this package (the ONLY place the two types meet), so the executor -// borrows it rather than re-implementing the edge. -func IssueToProto(si store.Issue) *compassv1.Issue { - return issueToProto(si) -} - // issueToProto maps the store-native issue to the canonical wire Issue. store // enums -> proto enums by value (they mirror: IssueState 0..8, ForgeProvider // 0..3); Forge is rebuilt as &compassv1.ForgeRef{Provider, Host}; Labels copied; -// empty->nil per the module contract; tracker/prs left nil (their producing -// slices own them). AgentAttribution is set only for a Compass-authored issue +// empty->nil per the module contract; tracker/prs left nil (prs are loaded by +// record caches wire under its id and coordinate. Callers hold p.mu. +func (p *IssueProjection) record(wire *compassv1.Issue, coord store.ForgeCoord) { + p.issues[wire.GetId()] = wire + p.byCoord[coord] = wire.GetId() +} + +// IssueToProtoWithPrs). AgentAttribution is set only for a Compass-authored issue // (a non-empty agent_handle); a human author leaves it unset. func issueToProto(si store.Issue) *compassv1.Issue { out := &compassv1.Issue{ diff --git a/go/internal/board/issue_projection_test.go b/go/internal/board/issue_projection_test.go index 825c12478..0c05ae34a 100644 --- a/go/internal/board/issue_projection_test.go +++ b/go/internal/board/issue_projection_test.go @@ -8,6 +8,7 @@ package board // PublishIssueUpdate/Rehydrate hit the store and live in the pgtest suite. import ( + "context" "testing" "time" @@ -207,7 +208,9 @@ func TestRecordAndPublishFansAndRecordsCommittedState(t *testing.T) { Number: 42, State: store.IssueStateInProgress, } - p.RecordAndPublish(committed) + if err := p.RecordAndPublish(context.Background(), committed); err != nil { + t.Fatalf("RecordAndPublish: %v", err) + } select { case e, ok := <-sub.Live: diff --git a/go/internal/board/issue_prs.go b/go/internal/board/issue_prs.go new file mode 100644 index 000000000..086ef8f31 --- /dev/null +++ b/go/internal/board/issue_prs.go @@ -0,0 +1,173 @@ +//go:build unix + +package board + +import ( + "context" + "fmt" + "slices" + + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/ingest" + "github.com/RigelBuild/compass/go/internal/store" + "google.golang.org/protobuf/encoding/protojson" +) + +// PublishPullRequestUpdate stores a hydrated PR and its closing references, then +// republishes every board issue linked to it before or after the write with a +// fresh prs list. Only Prs changes on the cached issue, so a concurrent state +// change is never rolled back. +func (p *IssueProjection) PublishPullRequestUpdate(ctx context.Context, in ingest.IngestedPullRequest) error { + row, refs, err := pullRequestRow(in) + if err != nil { + return err + } + affected, err := p.store.UpsertPullRequest(ctx, row, refs) + if err != nil { + return fmt.Errorf("board: upsert pull request: %w", err) + } + return p.republishPrs(ctx, affected) +} + +// pullRequestRow maps an ingested PR to the store row plus its closing refs on the PR's forge. +func pullRequestRow(in ingest.IngestedPullRequest) (store.PullRequestRow, []store.ForgeCoord, error) { + pr := in.PR + body, err := protojson.Marshal(pr) + if err != nil { + return store.PullRequestRow{}, nil, fmt.Errorf("board: marshal pull request: %w", err) + } + provider := store.ForgeProvider(pr.GetForge().GetProvider()) + host := pr.GetForge().GetHost() + row := store.PullRequestRow{ + Coord: store.ForgeCoord{Provider: provider, Host: host, Repo: pr.GetRepo(), Number: uint64(pr.GetNumber())}, + State: pr.GetForgeState(), + CreatedAt: in.CreatedAt, + UpdatedAt: in.UpdatedAt, + PR: body, + } + refs := make([]store.ForgeCoord, 0, len(in.ClosingRefs)) + for _, r := range in.ClosingRefs { + refs = append(refs, store.ForgeCoord{Provider: provider, Host: host, Repo: r.Repo, Number: r.Number}) + } + return row, refs, nil +} + +// storeIssueCoord is issueCoord for a store row. +func storeIssueCoord(si store.Issue) store.ForgeCoord { + return store.ForgeCoord{Provider: si.ForgeProvider, Host: si.ForgeHost, Repo: si.Repo, Number: uint64(si.Number)}.Normalized() +} + +// loadPrs reads the ordered wire prs for each coordinate in one query. +func (p *IssueProjection) loadPrs(ctx context.Context, coords []store.ForgeCoord) (map[store.ForgeCoord][]*compassv1.PullRequest, error) { + // A number-0 row is not a forge coordinate and can carry no links. + coords = slices.DeleteFunc(slices.Clone(coords), func(c store.ForgeCoord) bool { return c.Number == 0 }) + if len(coords) == 0 { + return map[store.ForgeCoord][]*compassv1.PullRequest{}, nil + } + rows, err := p.store.PullRequestsForIssues(ctx, coords) + if err != nil { + return nil, fmt.Errorf("board: load pull requests: %w", err) + } + out := make(map[store.ForgeCoord][]*compassv1.PullRequest, len(rows)) + for c, prs := range rows { + wire := make([]*compassv1.PullRequest, 0, len(prs)) + for _, r := range prs { + pr := &compassv1.PullRequest{} + // A row written by a newer binary must not fail a rollback's startup. + if err := (protojson.UnmarshalOptions{DiscardUnknown: true}).Unmarshal(r.PR, pr); err != nil { + return nil, fmt.Errorf("board: decode stored pull request %v: %w", r.Coord, err) + } + wire = append(wire, pr) + } + out[c] = wire + } + return out, nil +} + +// republishPrs reloads prs for the given issues and publishes each one already +// on the board. Issues not on the board are skipped; their links wait. Each +// caller reloads after its own commit, so the last to take writeMu applies the +// newest list. +func (p *IssueProjection) republishPrs(ctx context.Context, coords []store.ForgeCoord) error { + p.writeMu.Lock() + defer p.writeMu.Unlock() + return p.republishPrsLocked(ctx, coords) +} + +// republishPrsLocked is republishPrs for a caller holding writeMu. +func (p *IssueProjection) republishPrsLocked(ctx context.Context, coords []store.ForgeCoord) error { + if len(coords) == 0 { + return nil + } + prs, err := p.loadPrs(ctx, coords) + if err != nil { + return err + } + p.mu.Lock() + defer p.mu.Unlock() + for _, c := range coords { + id, ok := p.byCoord[c.Normalized()] + if !ok { + continue + } + next := cloneIssue(p.issues[id]) + next.Prs = prs[c.Normalized()] + p.issues[id] = next + p.bus.Publish(&compassv1.SubscribeEventsResponse{ + Payload: &compassv1.SubscribeEventsResponse_Issue{Issue: next}, + }) + } + return nil +} + +// CommittedIssue maps a committed row to the wire Issue, taking prs from the +// cache; a row not yet cached loads them from the store. +func (p *IssueProjection) CommittedIssue(ctx context.Context, si store.Issue) (*compassv1.Issue, error) { + if p != nil { + p.mu.RLock() + prev, ok := p.issues[si.ID] + var prs []*compassv1.PullRequest + if ok { + prs = cloneIssue(prev).GetPrs() + } + p.mu.RUnlock() + if ok { + wire := issueToProto(si) + wire.Prs = prs + return wire, nil + } + } + wires, err := p.IssueToProtoWithPrs(ctx, si) + if err != nil { + return nil, err + } + return wires[0], nil +} + +// IssueToProtoWithPrs maps committed rows to wire Issues with their prs loaded, +// for responses built outside the cache (SetIssueState, SearchIssues). A nil +// projection, as in tests that run without a board, maps rows without prs. +func (p *IssueProjection) IssueToProtoWithPrs(ctx context.Context, rows ...store.Issue) ([]*compassv1.Issue, error) { + if p == nil { + out := make([]*compassv1.Issue, 0, len(rows)) + for _, si := range rows { + out = append(out, issueToProto(si)) + } + return out, nil + } + coords := make([]store.ForgeCoord, 0, len(rows)) + for _, si := range rows { + coords = append(coords, storeIssueCoord(si)) + } + prs, err := p.loadPrs(ctx, coords) + if err != nil { + return nil, err + } + out := make([]*compassv1.Issue, 0, len(rows)) + for i, si := range rows { + wire := issueToProto(si) + wire.Prs = prs[coords[i]] + out = append(out, wire) + } + return out, nil +} diff --git a/go/internal/board/issue_prs_pgtest_test.go b/go/internal/board/issue_prs_pgtest_test.go new file mode 100644 index 000000000..330d6efe9 --- /dev/null +++ b/go/internal/board/issue_prs_pgtest_test.go @@ -0,0 +1,290 @@ +//go:build pgtest && unix + +package board + +// The projection keeps each board issue's prs current: on PR hydrate, on issue +// upsert, on state change and across Rehydrate. + +import ( + "context" + "slices" + "testing" + "time" + + "github.com/RigelBuild/compass/go/events" + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/forge" + "github.com/RigelBuild/compass/go/internal/ingest" + "github.com/RigelBuild/compass/go/internal/store" +) + +var prT0 = time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + +// ingestedPR is a hydrated PR in RigelBuild/compass closing the given issue numbers. +func ingestedPR(number uint32, age, updated time.Duration, closes ...uint64) ingest.IngestedPullRequest { + refs := make([]forge.IssueRef, 0, len(closes)) + for _, n := range closes { + refs = append(refs, forge.IssueRef{Repo: "RigelBuild/Compass", Number: n}) + } + return ingest.IngestedPullRequest{ + PR: &compassv1.PullRequest{ + Forge: &compassv1.ForgeRef{Provider: compassv1.ForgeProvider_FORGE_PROVIDER_GITHUB, Host: "github.com"}, + Repo: "RigelBuild/compass", + Number: number, + ForgeState: "open", + }, + CreatedAt: prT0.Add(age), + UpdatedAt: prT0.Add(updated), + ClosingRefs: refs, + } +} + +func prNumbers(iss *compassv1.Issue) []uint32 { + out := make([]uint32, 0, len(iss.GetPrs())) + for _, pr := range iss.GetPrs() { + out = append(out, pr.GetNumber()) + } + return out +} + +// cached returns the board's cached issue at number. +func cached(t *testing.T, p *IssueProjection, number uint32) *compassv1.Issue { + t.Helper() + for _, iss := range p.Snapshot() { + if iss.GetNumber() == number { + return iss + } + } + t.Fatalf("issue #%d not on the board", number) + return nil +} + +func subscribe(t *testing.T, bus *events.Bus[busPayload]) <-chan events.Stamped[busPayload] { + t.Helper() + sub, err := bus.Subscribe(0, bus.InstanceEpoch()) + if err != nil { + t.Fatalf("Subscribe: %v", err) + } + t.Cleanup(sub.Cancel) + return sub.Live +} + +func mustPublishIssue(t *testing.T, p *IssueProjection, number uint32) { + t.Helper() + if err := p.PublishIssueUpdate(context.Background(), canonicalIssue(number)); err != nil { + t.Fatalf("PublishIssueUpdate #%d: %v", number, err) + } +} + +func mustPublishPR(t *testing.T, p *IssueProjection, in ingest.IngestedPullRequest) { + t.Helper() + if err := p.PublishPullRequestUpdate(context.Background(), in); err != nil { + t.Fatalf("PublishPullRequestUpdate: %v", err) + } +} + +func TestPublishPullRequestAttachesAndStateChangeKeepsPrs(t *testing.T) { + ctx := context.Background() + p, bus, st := newIssueBoard(t) + mustPublishIssue(t, p, 1) + live := subscribe(t, bus) + + mustPublishPR(t, p, ingestedPR(10, 0, time.Hour, 1)) + if got := prNumbers(recvIssue(t, live)); !slices.Equal(got, []uint32{10}) { + t.Fatalf("fanned prs = %v, want [10]", got) + } + + id := cached(t, p, 1).GetId() + if err := st.SetIssueState(ctx, id, store.IssueStateInProgress); err != nil { + t.Fatalf("SetIssueState: %v", err) + } + committed, err := st.GetIssue(ctx, id) + if err != nil { + t.Fatalf("GetIssue: %v", err) + } + if err := p.RecordAndPublish(context.Background(), committed); err != nil { + t.Fatalf("RecordAndPublish: %v", err) + } + got := recvIssue(t, live) + if got.GetState() != compassv1.IssueState_ISSUE_STATE_IN_PROGRESS || !slices.Equal(prNumbers(got), []uint32{10}) { + t.Fatalf("state change fanned state %v prs %v, want IN_PROGRESS [10]", got.GetState(), prNumbers(got)) + } + wire, err := p.IssueToProtoWithPrs(ctx, committed) + if err != nil || !slices.Equal(prNumbers(wire[0]), []uint32{10}) { + t.Fatalf("IssueToProtoWithPrs = %v, %v; want prs [10]", wire, err) + } +} + +func TestPublishPullRequestRepublishesIssueThatLostRef(t *testing.T) { + p, bus, _ := newIssueBoard(t) + mustPublishIssue(t, p, 1) + mustPublishIssue(t, p, 2) + mustPublishPR(t, p, ingestedPR(10, 0, time.Hour, 1, 2)) + + live := subscribe(t, bus) + mustPublishPR(t, p, ingestedPR(10, 0, 2*time.Hour, 2)) + fanned := map[uint32][]uint32{} + for range 2 { + iss := recvIssue(t, live) + fanned[iss.GetNumber()] = prNumbers(iss) + } + if len(fanned[1]) != 0 || !slices.Equal(fanned[2], []uint32{10}) { + t.Fatalf("fanned = %v, want #1 without the PR and #2 with it", fanned) + } +} + +func TestPullRequestBeforeIssueAttachesLater(t *testing.T) { + p, _, _ := newIssueBoard(t) + mustPublishPR(t, p, ingestedPR(10, 0, time.Hour, 1)) + mustPublishIssue(t, p, 1) + if got := prNumbers(cached(t, p, 1)); !slices.Equal(got, []uint32{10}) { + t.Fatalf("late issue prs = %v, want [10]", got) + } +} + +func TestFallbackMovesWhenExplicitTargetArrives(t *testing.T) { + ctx := context.Background() + p, bus, st := newIssueBoard(t) + agent := mustBoardAgent(t, st) + mustPublishIssue(t, p, 1) // issue A, the closing ref + + // PR 10 explicitly targets B (#2, not yet on the board) and closes A. + created := ingestedPR(10, 0, 0, 1) + row, _, err := pullRequestRow(created) + if err != nil { + t.Fatalf("pullRequestRow: %v", err) + } + row.UpdatedAt = time.Unix(0, 0) + target := store.ForgeCoord{Provider: store.ForgeProviderGitHub, Host: "github.com", Repo: "rigelbuild/compass", Number: 2} + authored := store.AuthoredArtifact{ + Provider: row.Coord.Provider, Host: row.Coord.Host, Repo: "RigelBuild/compass", + Kind: store.ForgeArtifactKindPullRequest, Number: 10, + AgentAccountID: agent.agent, OwnerUserID: agent.owner, SessionID: "s", CreatedAtUnixMS: 1, + } + if err := st.CreatePullRequestWithLink(ctx, authored, row, &target); err != nil { + t.Fatalf("CreatePullRequestWithLink: %v", err) + } + mustPublishPR(t, p, ingestedPR(10, 0, time.Hour, 1)) + if got := prNumbers(cached(t, p, 1)); !slices.Equal(got, []uint32{10}) { + t.Fatalf("A prs before B arrives = %v, want the fallback [10]", got) + } + + live := subscribe(t, bus) + mustPublishIssue(t, p, 2) + fanned := map[uint32][]uint32{} + for range 2 { + iss := recvIssue(t, live) + fanned[iss.GetNumber()] = prNumbers(iss) + } + if !slices.Equal(fanned[2], []uint32{10}) || len(fanned[1]) != 0 { + t.Fatalf("fanned = %v, want B with [10] and A republished without it", fanned) + } +} + +func TestPrsOrderSurvivesRehydrate(t *testing.T) { + p, _, st := newIssueBoard(t) + mustPublishIssue(t, p, 1) + mustPublishPR(t, p, ingestedPR(30, 3*time.Hour, 3*time.Hour, 1)) + mustPublishPR(t, p, ingestedPR(10, time.Hour, 4*time.Hour, 1)) + mustPublishPR(t, p, ingestedPR(20, 2*time.Hour, 5*time.Hour, 1)) + + bus := events.NewBus[busPayload]() + t.Cleanup(bus.Close) + fresh := NewIssueProjection(bus, st) + if err := fresh.Rehydrate(context.Background()); err != nil { + t.Fatalf("Rehydrate: %v", err) + } + if got := prNumbers(cached(t, fresh, 1)); !slices.Equal(got, []uint32{10, 20, 30}) { + t.Fatalf("rehydrated prs = %v, want forge-created order [10 20 30]", got) + } +} + +type boardAgent struct{ agent, owner store.AccountID } + +func mustBoardAgent(t *testing.T, st *store.Store) boardAgent { + t.Helper() + ctx := context.Background() + u, err := st.CreateUser(ctx, store.NewUser{Handle: "prowner"}) + if err != nil { + t.Fatalf("CreateUser: %v", err) + } + a, err := st.CreateAgent(ctx, u.ID, store.NewAgent{Handle: "pragent"}) + if err != nil { + t.Fatalf("CreateAgent: %v", err) + } + return boardAgent{agent: a.ID, owner: u.ID} +} + +// TestTransitionOnUncachedIssueLoadsPrs pins the cache-miss path: a transition +// on an issue not yet in the cache fans and caches its stored prs. +func TestTransitionOnUncachedIssueLoadsPrs(t *testing.T) { + ctx := context.Background() + p, bus, st := newIssueBoard(t) + mustPublishIssue(t, p, 1) + mustPublishPR(t, p, ingestedPR(10, 0, time.Hour, 1)) + id := cached(t, p, 1).GetId() + + cold := NewIssueProjection(bus, st) + if err := st.SetIssueState(ctx, id, store.IssueStateInProgress); err != nil { + t.Fatalf("SetIssueState: %v", err) + } + committed, err := st.GetIssue(ctx, id) + if err != nil { + t.Fatalf("GetIssue: %v", err) + } + if err := cold.RecordAndPublish(ctx, committed); err != nil { + t.Fatalf("RecordAndPublish: %v", err) + } + wire, err := cold.CommittedIssue(ctx, committed) + if err != nil || !slices.Equal(prNumbers(wire), []uint32{10}) { + t.Fatalf("CommittedIssue prs = %v, %v; want [10]", prNumbers(wire), err) + } +} + +// TestTransitionRecordsRowNewerThanCaller pins the re-read: a forge update that +// commits after the transition's read is not rolled back in the cache. +func TestTransitionRecordsRowNewerThanCaller(t *testing.T) { + ctx := context.Background() + p, _, st := newIssueBoard(t) + mustPublishIssue(t, p, 1) + id := cached(t, p, 1).GetId() + if err := st.SetIssueState(ctx, id, store.IssueStateInProgress); err != nil { + t.Fatalf("SetIssueState: %v", err) + } + stale, err := st.GetIssue(ctx, id) + if err != nil { + t.Fatalf("GetIssue: %v", err) + } + newer := canonicalIssue(1) + newer.Title = "renamed" + if err := p.PublishIssueUpdate(ctx, newer); err != nil { + t.Fatalf("PublishIssueUpdate: %v", err) + } + if err := p.RecordAndPublish(ctx, stale); err != nil { + t.Fatalf("RecordAndPublish: %v", err) + } + if got := cached(t, p, 1); got.GetTitle() != "renamed" || got.GetState() != compassv1.IssueState_ISSUE_STATE_IN_PROGRESS { + t.Fatalf("cached title %q state %v, want renamed IN_PROGRESS", got.GetTitle(), got.GetState()) + } +} + +// TestTransitionPublishesWhenReadBackFails pins that a durable transition still +// reaches the board when the re-read fails, and the error is reported. +func TestTransitionPublishesWhenReadBackFails(t *testing.T) { + p, _, st := newIssueBoard(t) + mustPublishIssue(t, p, 1) + id := cached(t, p, 1).GetId() + committed, err := st.GetIssue(context.Background(), id) + if err != nil { + t.Fatalf("GetIssue: %v", err) + } + committed.State = store.IssueStateInProgress + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := p.RecordAndPublish(ctx, committed); err == nil { + t.Fatal("RecordAndPublish error = nil, want the failed re-read") + } + if got := cached(t, p, 1).GetState(); got != compassv1.IssueState_ISSUE_STATE_IN_PROGRESS { + t.Fatalf("cached state = %v, want IN_PROGRESS", got) + } +} diff --git a/go/internal/ingest/pull_request_sink.go b/go/internal/ingest/pull_request_sink.go new file mode 100644 index 000000000..396f8f0d8 --- /dev/null +++ b/go/internal/ingest/pull_request_sink.go @@ -0,0 +1,16 @@ +package ingest + +import ( + "time" + + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/forge" +) + +// IngestedPullRequest carries what the wire PullRequest lacks: forge times and closing refs. +type IngestedPullRequest struct { + PR *compassv1.PullRequest + CreatedAt time.Time + UpdatedAt time.Time + ClosingRefs []forge.IssueRef +} diff --git a/go/internal/store/db/pull_requests.sql.go b/go/internal/store/db/pull_requests.sql.go index 6f41abb2d..be945d5de 100644 --- a/go/internal/store/db/pull_requests.sql.go +++ b/go/internal/store/db/pull_requests.sql.go @@ -268,8 +268,9 @@ SELECT l.issue_forge_provider, l.issue_forge_host, l.issue_repo, l.issue_number, FROM pull_request_issue_links e JOIN issues i ON i.tenant_id = e.tenant_id AND i.forge_provider = e.issue_forge_provider - AND i.forge_host = e.issue_forge_host AND i.repo = e.issue_repo - AND i.number = e.issue_number + AND i.forge_host = e.issue_forge_host AND i.number = e.issue_number + -- issues keeps the ingested repo casing; links hold GitHub repos lowercased. + AND (CASE WHEN i.forge_provider = 1 THEN lower(i.repo) ELSE i.repo END) = e.issue_repo WHERE e.source = 1 AND e.tenant_id = l.tenant_id AND e.pr_forge_provider = l.pr_forge_provider AND e.pr_forge_host = l.pr_forge_host AND e.pr_repo = l.pr_repo AND e.pr_number = l.pr_number) diff --git a/go/internal/store/migrations/0001_init.sql b/go/internal/store/migrations/0001_init.sql index ad7af0541..ac81779b9 100644 --- a/go/internal/store/migrations/0001_init.sql +++ b/go/internal/store/migrations/0001_init.sql @@ -957,6 +957,11 @@ CREATE TABLE issues ( CREATE UNIQUE INDEX issues_coordinate_key ON issues (tenant_id, forge_provider, forge_host, repo, number); +-- The board join matches GitHub issue repos case-insensitively. +CREATE INDEX issues_board_coordinate_idx ON issues + (tenant_id, forge_provider, forge_host, + (CASE WHEN forge_provider = 1 THEN lower(repo) ELSE repo END), number); + CREATE INDEX issues_search_idx ON issues USING gin (search_tsv); -- ── Forge subscriptions & reconcile watermarks ─────────────────────────────── diff --git a/go/internal/store/pull_requests.go b/go/internal/store/pull_requests.go index cbbdf0f6f..3fa04eeea 100644 --- a/go/internal/store/pull_requests.go +++ b/go/internal/store/pull_requests.go @@ -34,8 +34,8 @@ const ( linkSourceClosingRef int16 = 2 ) -// normalized trims and, on GitHub, lowercases the repo so board joins never miss on case. -func (c ForgeCoord) normalized() ForgeCoord { +// Normalized trims and, on GitHub, lowercases the repo so board joins never miss on case. +func (c ForgeCoord) Normalized() ForgeCoord { c.Repo = strings.TrimSpace(c.Repo) if c.Provider == ForgeProviderGitHub { c.Repo = strings.ToLower(c.Repo) @@ -80,7 +80,7 @@ func coordFromDB(provider int16, host, repo string, number int64) ForgeCoord { } func (pr PullRequestRow) normalized() (PullRequestRow, error) { - pr.Coord = pr.Coord.normalized() + pr.Coord = pr.Coord.Normalized() if err := pr.Coord.valid(); err != nil { return PullRequestRow{}, err } @@ -98,7 +98,7 @@ func normalizeCoords(in []ForgeCoord) ([]ForgeCoord, error) { out := make([]ForgeCoord, 0, len(in)) seen := make(map[ForgeCoord]bool, len(in)) for _, c := range in { - c = c.normalized() + c = c.Normalized() if err := c.valid(); err != nil { return nil, err } @@ -125,13 +125,13 @@ func (s *Store) CreatePullRequestWithLink(ctx context.Context, a AuthoredArtifac if err := a.valid(); err != nil { return err } - authored := ForgeCoord{Provider: a.Provider, Host: a.Host, Repo: a.Repo, Number: a.Number}.normalized() + authored := ForgeCoord{Provider: a.Provider, Host: a.Host, Repo: a.Repo, Number: a.Number}.Normalized() if a.Kind != ForgeArtifactKindPullRequest || authored != pr.Coord { return fmt.Errorf("%w: authored artifact does not name the pull request", ErrInvalidArgument) } var target ForgeCoord if issue != nil { - target = issue.normalized() + target = issue.Normalized() if err := target.valid(); err != nil { return err } @@ -260,8 +260,8 @@ func replaceClosingRefs(ctx context.Context, qtx *db.Queries, p forgeCoordDB, re } // PullRequestsForIssues returns the PRs attached to each issue, oldest forge -// creation first, keyed by the normalized issue coordinate. A closing reference attaches a PR only when none of its -// explicit targets is a board issue. +// creation first, keyed by the normalized issue coordinate. A closing reference +// attaches a PR only when none of its explicit targets is a board issue. func (s *Store) PullRequestsForIssues(ctx context.Context, issues []ForgeCoord) (map[ForgeCoord][]PullRequestRow, error) { want, err := normalizeCoords(issues) if err != nil { @@ -305,7 +305,7 @@ func (s *Store) PullRequestsForIssues(ctx context.Context, issues []ForgeCoord) // explicit target is issue: the issues that gain or lose the PR when issue // enters or leaves the board. func (s *Store) FallbackIssuesForTarget(ctx context.Context, issue ForgeCoord) ([]ForgeCoord, error) { - issue = issue.normalized() + issue = issue.Normalized() if err := issue.valid(); err != nil { return nil, err } @@ -326,7 +326,7 @@ func (s *Store) FallbackIssuesForTarget(ctx context.Context, issue ForgeCoord) ( // 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() + pr = pr.Normalized() if err := pr.valid(); err != nil { return time.Time{}, false, err } diff --git a/go/internal/store/pull_requests_pgtest_test.go b/go/internal/store/pull_requests_pgtest_test.go index 5e5cb7da7..4c1e633cf 100644 --- a/go/internal/store/pull_requests_pgtest_test.go +++ b/go/internal/store/pull_requests_pgtest_test.go @@ -335,3 +335,24 @@ func TestPRsBackfilledMark(t *testing.T) { t.Fatalf("unknown repo err = %v, want ErrNotFound", err) } } + +// Board issues keep their ingested repo casing; a mixed-case board row still +// counts as a resolvable explicit target. +func TestPullRequestMixedCaseBoardIssueResolvesExplicitLink(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + agent, owner := seedAgent(t, s, "prl11") + seedBoardIssue(t, ctx, s, ghCoord("Owner/Repo", 1)) + closing := ghCoord("owner/repo", 2) + seedBoardIssue(t, ctx, s, closing) + + target := ghCoord("owner/repo", 1) + pr := prRow("owner/repo", 10, "open", prBase, time.Unix(0, 0)) + if err := s.CreatePullRequestWithLink(ctx, authoredPR(agent, owner, pr), pr, &target); err != nil { + t.Fatalf("CreatePullRequestWithLink: %v", err) + } + mustUpsertPR(t, ctx, s, prRow("owner/repo", 10, "open", prBase, prBase.Add(time.Hour)), closing) + if got := prsFor(t, ctx, s, closing); len(got) != 0 { + t.Fatalf("closing ref attached although the explicit target is on the board: %v", got) + } +} diff --git a/go/internal/store/queries/pull_requests.sql b/go/internal/store/queries/pull_requests.sql index 5b57305ab..9c381a799 100644 --- a/go/internal/store/queries/pull_requests.sql +++ b/go/internal/store/queries/pull_requests.sql @@ -74,8 +74,9 @@ SELECT l.issue_forge_provider, l.issue_forge_host, l.issue_repo, l.issue_number, FROM pull_request_issue_links e JOIN issues i ON i.tenant_id = e.tenant_id AND i.forge_provider = e.issue_forge_provider - AND i.forge_host = e.issue_forge_host AND i.repo = e.issue_repo - AND i.number = e.issue_number + AND i.forge_host = e.issue_forge_host AND i.number = e.issue_number + -- issues keeps the ingested repo casing; links hold GitHub repos lowercased. + AND (CASE WHEN i.forge_provider = 1 THEN lower(i.repo) ELSE i.repo END) = e.issue_repo WHERE e.source = 1 AND e.tenant_id = l.tenant_id AND e.pr_forge_provider = l.pr_forge_provider AND e.pr_forge_host = l.pr_forge_host AND e.pr_repo = l.pr_repo AND e.pr_number = l.pr_number) diff --git a/go/server/board.go b/go/server/board.go index aec834ab6..b0e006f68 100644 --- a/go/server/board.go +++ b/go/server/board.go @@ -10,6 +10,7 @@ import ( "context" "errors" "fmt" + "log/slog" "sync" "connectrpc.com/connect" @@ -119,7 +120,11 @@ func (b *boardService) SetIssueStateAsAccount( if err != nil { return nil, err } - return &compassv1internal.SetIssueStateResponse{Issue: board.IssueToProto(committed)}, nil + wire, err := b.issueBrd.CommittedIssue(ctx, committed) + if err != nil { + return nil, connect.NewError(connect.CodeInternal, err) + } + return &compassv1internal.SetIssueStateResponse{Issue: wire}, nil } // errUnspecifiedTarget is the in-band cause for an ISSUE_STATE_UNSPECIFIED @@ -180,11 +185,11 @@ func (b *boardService) SetIssueState( return store.Issue{}, transitionStoreError(err) } - // Record + fan out the committed transition on the projection (the issue=16 - // live stream + the durable-cache snapshot). State-only: the executor already - // owns the durable commit, so this never touches the store (a forge-only - // upsert would demand forge fields and could not carry the state column). - b.issueBrd.RecordAndPublish(committed) + // Record + fan out the committed transition. It always publishes; an error + // only means its re-read failed, and the transition is already durable. + if err := b.issueBrd.RecordAndPublish(ctx, committed); err != nil { + slog.WarnContext(ctx, "board: issue transition published from the caller's row", "issue", issueID, "error", err) + } // Outbound tracker mirror on a real transition. Nil-safe. ARCHIVED has no // tracker status, so it is elided. The mirror runs AFTER the state is durable + diff --git a/go/server/service.go b/go/server/service.go index c02bfd1e5..92dd616fa 100644 --- a/go/server/service.go +++ b/go/server/service.go @@ -314,9 +314,9 @@ func (s *service) SearchIssues( } return nil, connect.NewError(connect.CodeInternal, err) } - out := make([]*compassv1.Issue, 0, len(issues)) - for _, issue := range issues { - out = append(out, board.IssueToProto(issue)) + out, err := s.issueBrd.IssueToProtoWithPrs(ctx, issues...) + if err != nil { + return nil, connect.NewError(connect.CodeInternal, err) } return connect.NewResponse(&compassv1.SearchIssuesResponse{Issues: out}), nil }