package pipeline import ( "context" "fmt" "net/http" "os" "strconv" "strings" "sync" "evobgp/internal/store" ) func collectConcurrency() int { n := 8 if s := strings.TrimSpace(os.Getenv("EVOBGP_COLLECT_CONCURRENCY")); s != "" { if v, err := strconv.Atoi(s); err == nil && v > 0 { n = v } } if n > 32 { n = 32 } return n } // aggregateTenantPrefixRowsAll builds the union of materialized prefixes for all enabled modules. // Uses per-module snapshots when inputs are unchanged to avoid duplicate external fetches on render. func aggregateTenantPrefixRowsAll(ctx context.Context, st store.Backend, hc *http.Client, tenantID string) ([]store.PrefixRow, error) { mods := st.ListModules(tenantID) type modRef struct { id string } var enabled []modRef for _, m := range mods { if m != nil && m.Enabled { enabled = append(enabled, modRef{id: m.ID}) } } if len(enabled) == 0 { return nil, nil } if len(enabled) == 1 { omod, err := st.GetModule(tenantID, enabled[0].id) if err != nil { return nil, err } return rowsForModule(ctx, st, hc, tenantID, omod) } sem := make(chan struct{}, collectConcurrency()) results := make([][]store.PrefixRow, len(enabled)) errs := make([]error, len(enabled)) var wg sync.WaitGroup for i, ref := range enabled { wg.Add(1) go func(idx int, moduleID string) { defer wg.Done() sem <- struct{}{} defer func() { <-sem }() omod, err := st.GetModule(tenantID, moduleID) if err != nil { errs[idx] = err return } rows, err := rowsForModule(ctx, st, hc, tenantID, omod) if err != nil { errs[idx] = fmt.Errorf("module %s: %w", moduleID, err) return } results[idx] = rows }(i, ref.id) } wg.Wait() for _, err := range errs { if err != nil { return nil, err } } var out []store.PrefixRow for _, part := range results { out = append(out, part...) } return out, nil } func rowsForModule(ctx context.Context, st store.Backend, hc *http.Client, tenantID string, mod *store.Module) ([]store.PrefixRow, error) { if rows, ok, err := moduleRowsFromSnapshot(st, tenantID, mod); err != nil { return nil, err } else if ok { return rows, nil } var prior []store.PrefixRow if snap, ok, _ := st.GetModulePrefixSnapshot(tenantID, mod.ID); ok && snap != nil { prior = snap.Prefixes } rows, err := collectModulePrefixRows(ctx, st, hc, tenantID, mod, prior) if err != nil { return nil, err } if err := persistModuleSnapshot(st, tenantID, mod, rows); err != nil { return nil, err } return rows, nil }