Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
151 changes: 91 additions & 60 deletions go/internal/board/issue_projection.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ package board

import (
"context"
"errors"
"fmt"
"sort"
"sync"
Expand All @@ -28,18 +29,24 @@ 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
// fans issue upserts onto and the store it reads/writes durable issue state
// 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),
}
}

Expand All @@ -48,30 +55,23 @@ 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).
id, err := p.store.UpsertIssueForgeFields(ctx, protoToForgeFields(issue))
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
Expand All @@ -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
Expand All @@ -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{
Expand Down
5 changes: 4 additions & 1 deletion go/internal/board/issue_projection_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ package board
// PublishIssueUpdate/Rehydrate hit the store and live in the pgtest suite.

import (
"context"
"testing"
"time"

Expand Down Expand Up @@ -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:
Expand Down
Loading
Loading