Files
Denozordec ee364c8b6d
CI / changes (push) Successful in 7s
CI / openapi (push) Has been skipped
CI / go (push) Successful in 1m56s
CI / docker-web (deploy/docker/evobgp-web/Dockerfile, , evobgp-web) (push) Has been skipped
CI / docker-web (deploy/docker/evobgp-web/Dockerfile, evobgp-all, evobgp-web-all) (push) Has been skipped
CI / docker-bird (push) Has been skipped
CI / bird2 (push) Successful in 15s
CI / docker-go-prime (push) Successful in 26s
CI / docker-go (deploy/docker/evobgp-agent/Dockerfile, , evobgp-agent) (push) Successful in 1m0s
CI / docker-go (evobgp-all, 1, deploy/docker/gobinary/Dockerfile, , evobgp-all) (push) Successful in 2m15s
CI / docker-go (evobgp-api, 1, deploy/docker/gobinary/Dockerfile, , evobgp-api) (push) Successful in 1m31s
CI / docker-go (evobgp-deploy, 0, deploy/docker/gobinary/Dockerfile, , evobgp-deploy) (push) Successful in 1m32s
CI / docker-go (evobgp-ingest, 0, deploy/docker/gobinary/Dockerfile, , evobgp-ingest) (push) Successful in 1m40s
CI / docker-go (evobgp-node, 0, deploy/docker/gobinary/Dockerfile, , evobgp-node) (push) Successful in 1m23s
CI / docker-go (evobgp-render, 0, deploy/docker/gobinary/Dockerfile, , evobgp-render) (push) Successful in 1m23s
CI / docker-go (evobgp-scheduler, 0, deploy/docker/gobinary/Dockerfile, , evobgp-scheduler) (push) Successful in 1m20s
feat: update collectModulePrefixRows to support prior snapshots and enhance CDN/domain prefix collection
Modified the collectModulePrefixRows function to accept an optional priorSnapshot parameter, allowing for more efficient data retrieval by skipping unnecessary CDN fetches. Refactored the logic for collecting prefix rows from AS, CDN, and domain sources to utilize dedicated functions, improving code organization and maintainability. Additionally, introduced caching for module prefix snapshots to optimize performance during refresh operations.
2026-05-19 10:22:33 +07:00

107 lines
2.5 KiB
Go

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
}