cli: add multi-cell fan-out helpers (groupReposByCell, fanOutCells) · Entire

cli: add multi-cell fan-out helpers (groupReposByCell, fanOutCells)

48d5ac5·

Soph·1w ago·4 files·+457 added/-0 removed

The data plane has no server-side cross-cell aggregator: a query over all of the caller's repos must be fanned out to each cell hosting them and merged client-side (the entire.io BFF's code-search pattern). PR #1616 inlines that orchestration into search_cmd.go; this extracts the generic layer so code search — and any later repo-set command — shares one implementation:

Also documents the three cell-routing shapes in CLAUDE.md.

Co-Authored-By: Claude Fable 5 noreply@anthropic.com

Sessions

Entire-API Cell Routing (which cell does a data-plane request go to?)

The data plane (entire-api) is deployed per jurisdiction; a repo placement lives in exactly one cell, user /me/* activity is consolidated in the caller's home cell, and no server-side cross-cell aggregator exists. The CLI therefore has exactly three routing shapes, mirroring the entire.io BFF:

Token rule: identity tokens are per-jurisdiction, not per-cell. Multi-cell callers must build one auth.CellClientFactory (NewEntireAPICellClientFactory) per operation — it resolves the login subject once and mints at most one token per jurisdiction. fanOutCells does this automatically; do not call NewEntireAPICellClient in a loop.

Session Strategy (cmd/entire/cli/strategy/)

The CLI uses a manual-commit strategy for managing session data and checkpoints.

package cli

import (
    "context"
    "sort"
    "strings"
    "sync"
    "time"

"github.com/entireio/cli/cmd/entire/cli/api"
    "github.com/entireio/cli/cmd/entire/cli/auth"
    "github.com/entireio/cli/cmd/entire/cli/logging"
    "github.com/entireio/cli/internal/coreapi"
)

// This file is the multi-cell counterpart to cell_target.go: where
// resolveRepoCellTarget routes ONE repo-scoped call to the cell hosting that
// repo, the helpers here route a query over ALL of the caller's repos — group
// the repo index by hosting cell, then ask each cell about its own repos and
// let the caller merge. That mirrors the entire.io BFF's fan-out
// (code-search.ts: index → group by cell → per-cell call → merge); no
// server-side aggregator exists, cells are strictly local.

// cellGroup is one entire-api cell plus the caller's repos hosted there — the
// unit of a multi-cell fan-out.
type cellGroup struct {
    // cell is the physical cell name (e.g. aws-eu-west-1), the grouping key —
    // each repo placement lives in exactly one cell. Empty when the index did
    // not report one (the group then routes by jurisdiction, or home).
    cell string
    // clusterSlug joins the group to the cluster catalog
    // (RepoIndexEntry.ClusterSlug ↔ Cluster.Slug) to resolve baseURL. The
    // catalog does not expose a cell field, so the slug — not the cell name —
    // is the only reliable join key. Several clusters may share a cell; any of
    // them reports the jurisdiction's apiUrl, so the first seen slug serves.
    clusterSlug string
    // jurisdiction is the lowercased jurisdiction label; it drives the identity
    // token audience and is the routing fallback when baseURL is empty.
    jurisdiction string
    // baseURL is the cell's resolved apiUrl (resolveCellBaseURLs). Empty means
    // "route by jurisdiction" — the auth layer then resolves the jurisdiction's
    // default cell from the catalog.
    baseURL string
    // repoIDs are the caller's repo ULIDs placed in this cell, so the cell is
    // only ever asked about repos it hosts.
    repoIDs []string
}

// groupReposByCell groups a repo index by hosting cell, one group per distinct
// cell, deterministically ordered by cell name. Entries without an ID are
// skipped (nothing to ask the cell about).
func groupReposByCell(repos []coreapi.RepoIndexEntry) []cellGroup {
    byCell := make(map[string]*cellGroup)
    for _, r := range repos {
        id := strings.TrimSpace(r.ID)
        if id == "" {
            continue
        }
        key := strings.ToLower(strings.TrimSpace(r.Cell))
        g, ok := byCell[key]
        if !ok {
            g = &cellGroup{
                cell:         key,
                clusterSlug:  strings.ToLower(strings.TrimSpace(r.ClusterSlug)),
                jurisdiction: strings.ToLower(strings.TrimSpace(r.Jurisdiction)),
            }
            byCell[key] = g
        }
        g.repoIDs = append(g.repoIDs, id)
    }
    cells := make([]cellGroup, 0, len(byCell))
    for _, g := range byCell {
        cells = append(cells, *g)
    }
    sort.Slice(cells, func(i, j int) bool { return cells[i].cell < cells[j].cell })
    return cells
}

// resolveCellBaseURLs fills each group's baseURL from the cluster catalog,
// joining on ClusterSlug ↔ Cluster.Slug. Best-effort: on a catalog error or a
// missing/incomplete cluster row the group keeps baseURL "" and falls back to
// jurisdiction routing — a degraded catalog must not sink the fan-out.
func resolveCellBaseURLs(ctx context.Context, c cellCoreClient, cells []cellGroup) {
    clusters, err := c.ListClusters(ctx)
    if err != nil {
        logging.Debug(ctx, "cell fan-out: list clusters failed, using jurisdiction routing", "error", err.Error())
        return
    }
    bySlug := make(map[string]coreapi.Cluster, len(clusters.Clusters))
    for _, cl := range clusters.Clusters {
        bySlug[strings.ToLower(strings.TrimSpace(cl.Slug))] = cl
    }
    for i := range cells {
        cl, ok := bySlug[cells[i].clusterSlug]
        if !ok {
            logging.Debug(ctx, "cell fan-out: cluster not in catalog, using jurisdiction routing",
                "cluster_slug", cells[i].clusterSlug, "cell", cells[i].cell)
            continue
        }
        cells[i].baseURL = strings.TrimRight(strings.TrimSpace(cl.ApiUrl.Or("")), "/")
        if j := strings.ToLower(strings.TrimSpace(cl.Jurisdiction)); j != "" {
            cells[i].jurisdiction = j
        }
    }
}

// cellTarget converts the group's routing coordinates into the auth layer's
// CellTarget: full target when the catalog resolved a baseURL,
// jurisdiction-only when it didn't, nil (home routing) when neither is known.
func (g cellGroup) cellTarget() *auth.CellTarget {
    switch {
    case g.baseURL != "":
        return &auth.CellTarget{BaseURL: g.baseURL, Jurisdiction: g.jurisdiction}
    case g.jurisdiction != "":
        return &auth.CellTarget{Jurisdiction: g.jurisdiction}
    default:
        return nil
    }
}

// label names the group in errors and logs: cell, else jurisdiction, else home.
func (g cellGroup) label() string {
    switch {
    case g.cell != "":
        return g.cell
    case g.jurisdiction != "":
        return g.jurisdiction
    default:
        return "home"
    }
}

// fakeCellClientBuilder hands out unauthenticated clients keyed by target and
// records what it was asked for.
type fakeCellClientBuilder struct {
    mu      sync.Mutex // fanOutCells calls ClientFor from one goroutine per cell
    targets []*auth.CellTarget
    err     error
}

func (f *fakeCellClientBuilder) ClientFor(_ context.Context, target *auth.CellTarget) (*api.Client, error) {
    f.mu.Lock()
    f.targets = append(f.targets, target)
    f.mu.Unlock()
    if f.err != nil {
        return nil, f.err
    }
    base := "https://home.api.example"
    if target != nil && target.BaseURL != "" {
        base = target.BaseURL
    }
    return api.NewClientWithBaseURL("test-token", base), nil
}

func TestGroupReposByCell(t *testing.T) {
    t.Parallel()

repos := []coreapi.RepoIndexEntry{
        {ID: "01B", Cell: "aws-us-east-2", ClusterSlug: "us-prod", Jurisdiction: "us"},
        {ID: "01C", Cell: "AWS-US-EAST-2", ClusterSlug: "us-prod", Jurisdiction: "US"}, // case-folds into same group
        {ID: "01A", Cell: euWestCell, ClusterSlug: "eu-prod", Jurisdiction: "eu"},
        {ID: "", Cell: euWestCell}, // no ID → skipped
        {ID: "01D"},                // no cell → its own "" group (home/jurisdiction routing)
    }

cells := groupReposByCell(repos)

if len(cells) != 3 {
        t.Fatalf("groups = %d, want 3: %+v", len(cells), cells)
    }

// Deterministic order by cell name: "" < aws-eu-west-1 < aws-us-east-2.
    if cells[0].cell != "" || cells[1].cell != euWestCell || cells[2].cell != "aws-us-east-2" {
        t.Fatalf("group order = [%q %q %q], want [\"\" aws-eu-west-1 aws-us-east-2]", cells[0].cell, cells[1].cell, cells[2].cell)
    }

us := cells[2]
    if got := strings.Join(us.repoIDs, ","); got != "01B,01C" {
        t.Fatalf("us repoIDs = %q, want 01B,01C", got)
    }
    if us.clusterSlug != "us-prod" || us.jurisdiction != "us" {
        t.Fatalf("us group coordinates = %+v, want us-prod/us", us)
    }
}

func TestResolveCellBaseURLs_JoinsOnClusterSlug(t *testing.T) {
    t.Parallel()

cells := []cellGroup{
        {cell: euWestCell, clusterSlug: "eu-prod", jurisdiction: "eu"},
        {cell: "aws-ap-south-1", clusterSlug: "ap-prod", jurisdiction: "ap"}, // not in catalog
    }

fake := &fakeCellCore{clusters: []coreapi.Cluster{
        {Slug: "EU-Prod", Jurisdiction: "EU", ApiUrl: coreapi.NewOptString("https://aws-eu-west-1.api.entire.io/")},
    }}

resolveCellBaseURLs(context.Background(), fake, cells)

if got := cells[0].baseURL; got != "https://aws-eu-west-1.api.entire.io" {
        t.Fatalf("eu baseURL = %q, want the catalog apiUrl (trimmed)", got)
    }
    if cells[0].jurisdiction != "eu" {
        t.Fatalf("eu jurisdiction = %q, want normalised eu", cells[0].jurisdiction)
    }
    if cells[1].baseURL != "" {
        t.Fatalf("ap baseURL = %q, want empty (jurisdiction fallback)", cells[1].baseURL)
    }
}

func TestResolveCellBaseURLs_CatalogErrorLeavesJurisdictionRouting(t *testing.T) {
    t.Parallel()

cells := []cellGroup{{cell: euWestCell, clusterSlug: "eu-prod", jurisdiction: "eu"}}

resolveCellBaseURLs(context.Background(), &fakeCellCore{clustersErr: errors.New("boom")}, cells)

if cells[0].baseURL != "" {
        t.Fatalf("baseURL = %q, want empty after catalog error", cells[0].baseURL)
    }
}

func TestCellGroupTargetAndLabel(t *testing.T) {
    t.Parallel()

full := cellGroup{cell: euWestCell, jurisdiction: "eu", baseURL: "https://aws-eu-west-1.api.entire.io"}
    if tgt := full.cellTarget(); tgt == nil || tgt.BaseURL != full.baseURL || tgt.Jurisdiction != "eu" {
        t.Fatalf("full target = %+v", tgt)
    }

jur := cellGroup{jurisdiction: "eu"}
    if tgt := jur.cellTarget(); tgt == nil || tgt.BaseURL != "" || tgt.Jurisdiction != "eu" {
        t.Fatalf("jurisdiction-only target = %+v", tgt)
    }

if tgt := (cellGroup{}).cellTarget(); tgt != nil {
        t.Fatalf("empty group target = %+v, want nil (home routing)", tgt)
    }

if got := full.label(); got != euWestCell {
        t.Fatalf("label = %q", got)
    }
    if got := jur.label(); got != "eu" {
        t.Fatalf("label = %q", got)
    }
    if got := (cellGroup{}).label(); got != "home" {
        t.Fatalf("label = %q, want home", got)
    }
}

func TestFanOutCells_PartialFailureIsPerCell(t *testing.T) {
    // Not parallel: swaps the package-level newCellClientBuilder seam.
    withFakeCellClientBuilder(t, &fakeCellClientBuilder{})

cells := []cellGroup{
        {cell: euWestCell, jurisdiction: "eu", baseURL: "https://eu.api.example", repoIDs: []string{"01A"}},
        {cell: "aws-us-east-2", jurisdiction: "us", baseURL: "https://us.api.example", repoIDs: []string{"01B"}},
    }

boom := errors.New("cell down")

results, err := fanOutCells(context.Background(), false, time.Second, cells,
        func(ctx context.Context, g cellGroup, _ *api.Client) (string, error) {
            if _, ok := ctx.Deadline(); !ok {
                t.Error("per-cell ctx has no deadline")
            }
            if g.cell == euWestCell {
                return "", boom
            }
            return "hits:" + strings.Join(g.repoIDs, ","), nil
        })

if err != nil {
        t.Fatalf("fanOutCells: %v", err)
    }
    if len(results) != 2 {
        t.Fatalf("results = %d, want 2", len(results))
    }

// Input order preserved; the eu failure is isolated in its slot.
    if !errors.Is(results[0].err, boom) || results[0].group.cell != euWestCell {
        t.Fatalf("results[0] = %+v, want eu failure", results[0])
    }
    if results[1].err != nil || results[1].value != "hits:01B" {
        t.Fatalf("results[1] = %+v, want us success", results[1])
    }
}

func TestFanOutCells_SingleCellRunsSerially(t *testing.T) {
    // Not parallel: swaps the package-level newCellClientBuilder seam.
    builder := &fakeCellClientBuilder{}
    withFakeCellClientBuilder(t, builder)

cells := []cellGroup{{jurisdiction: "eu", baseURL: "https://eu.api.example"}}
    results, err := fanOutCells(context.Background(), false, time.Second, cells,
        func(_ context.Context, _ cellGroup, _ *api.Client) (string, error) {
            return "ok", nil
        })

if err != nil || len(results) != 1 || results[0].err != nil || results[0].value != "ok" {
        t.Fatalf("results = %+v, err = %v", results, err)
    }
    if len(builder.targets) != 1 || builder.targets[0].BaseURL != "https://eu.api.example" {
        t.Fatalf("builder targets = %+v", builder.targets)
    }
}

func TestFanOutCells_EmptyAndFactoryError(t *testing.T) {
    // Not parallel: swaps the package-level newCellClientBuilder seam.
    results, err := fanOutCells(context.Background(), false, time.Second, nil,
        func(context.Context, cellGroup, *api.Client) (int, error) { return 0, nil })
    if results != nil || err != nil {
        t.Fatalf("empty fan-out = (%v, %v), want (nil, nil)", results, err)
    }

factoryErr := errors.New("not logged in")
    prev := newCellClientBuilder
    newCellClientBuilder = func(context.Context, bool) (cellClientBuilder, error) { return nil, factoryErr }
    t.Cleanup(func() { newCellClientBuilder = prev })
    if _, err := fanOutCells(context.Background(), false, time.Second, []cellGroup{{jurisdiction: "eu"}},
        func(context.Context, cellGroup, *api.Client) (int, error) { return 0, nil }); !errors.Is(err, factoryErr) {
        t.Fatalf("err = %v, want factory error", err)
    }
}

// TestFanOutCells_ClientPerCellFromOneBuilder asserts every cell's client
// comes from the single shared builder (one subject, per-jurisdiction token
// reuse lives behind it in auth.CellClientFactory).
func TestFanOutCells_ClientPerCellFromOneBuilder(t *testing.T) {
    // Not parallel: swaps the package-level newCellClientBuilder seam.
    builder := &fakeCellClientBuilder{}
    withFakeCellClientBuilder(t, builder)
    var cells []cellGroup
    for i := range 3 {
        cells = append(cells, cellGroup{
            cell:         fmt.Sprintf("cell-%d", i),
            jurisdiction: "eu",
            baseURL:      fmt.Sprintf("https://cell-%d.api.example", i),
        })
    }

results, err := fanOutCells(context.Background(), false, time.Second, cells,
        func(_ context.Context, g cellGroup, _ *api.Client) (string, error) { return g.cell, nil })
    if err != nil {
        t.Fatalf("fanOutCells: %v", err)
    }
    for i, r := range results {
        if r.err != nil || r.value != fmt.Sprintf("cell-%d", i) {
            t.Fatalf("results[%d] = %+v", i, r)
        }
    }
    if len(builder.targets) != 3 {
        t.Fatalf("builder asked for %d targets, want 3", len(builder.targets))
    }
}

const euCellAPIURL = "https://eu.api.entire.io"

const euWestCell = "aws-eu-west-1"