Files
EvoBGP/internal/pipeline/collect_parallel.go
T
Denozordec 7a2b015b12
CI / changes (push) Successful in 11s
CI / openapi (push) Has been skipped
CI / go (push) Successful in 59s
CI / docker-web (deploy/docker/evobgp-web/Dockerfile, , evobgp-web) (push) Has been skipped
CI / docker-bird (push) Has been skipped
CI / docker-web (deploy/docker/evobgp-web/Dockerfile, evobgp-all, evobgp-web-all) (push) Has been skipped
CI / bird2 (push) Successful in 43s
CI / docker-go (evobgp-all, 1, deploy/docker/gobinary/Dockerfile, , evobgp-all) (push) Has been skipped
CI / docker-go (evobgp-api, 1, deploy/docker/gobinary/Dockerfile, , evobgp-api) (push) Has been skipped
CI / docker-go (evobgp-deploy, 0, deploy/docker/gobinary/Dockerfile, , evobgp-deploy) (push) Has been skipped
CI / docker-go (evobgp-ingest, 0, deploy/docker/gobinary/Dockerfile, , evobgp-ingest) (push) Has been skipped
CI / docker-go (evobgp-node, 0, deploy/docker/gobinary/Dockerfile, , evobgp-node) (push) Has been skipped
CI / docker-go (evobgp-render, 0, deploy/docker/gobinary/Dockerfile, , evobgp-render) (push) Has been skipped
CI / docker-go (evobgp-scheduler, 0, deploy/docker/gobinary/Dockerfile, , evobgp-scheduler) (push) Has been skipped
CI / docker-go-prime (push) Successful in 55s
CI / docker-go (deploy/docker/evobgp-agent/Dockerfile, , evobgp-agent) (push) Has been cancelled
feat: add tenant refresh functionality to HTTP API and job processing
Implemented a new endpoint for tenant refresh in the HTTP API, allowing for the refresh of modules associated with a tenant. Enhanced job processing to handle tenant refresh jobs, including logic for managing module IDs and job status updates. Updated the scheduler to enqueue tenant refresh jobs based on module due dates, improving the overall efficiency of module management. Additionally, introduced caching for ASN prefix data to optimize performance during refresh operations.
2026-05-19 10:35:20 +07:00

245 lines
5.8 KiB
Go

package pipeline
import (
"context"
"fmt"
"net/http"
"os"
"strings"
"sync"
"time"
"evobgp/internal/store"
)
func prefixRowsForSource(rows []store.PrefixRow, sourceKey string) []store.PrefixRow {
if len(rows) == 0 {
return nil
}
var out []store.PrefixRow
for _, row := range rows {
if row.Source == sourceKey {
out = append(out, row)
}
}
return out
}
func collectASPrefixRows(ctx context.Context, st store.Backend, hc *http.Client, tenantID string, mod *store.Module, list []*store.ASEntry) ([]store.PrefixRow, error) {
moduleID := mod.ID
legacy := strings.TrimSpace(os.Getenv("EVOBGP_ASN_RESOLVE")) == "0"
if legacy {
var rows []store.PrefixRow
for _, e := range list {
if !store.ValidASN(e.ASN) {
continue
}
comm := e.CommunityID
if comm == nil && mod.DefaultCommunityID != nil {
c := *mod.DefaultCommunityID
comm = &c
}
rows = append(rows, store.PrefixRow{Prefix: MaterializedASPrefixKey(e.ASN), CommunityID: comm, Source: "as_entry"})
}
return rows, nil
}
type entryResult struct {
rows []store.PrefixRow
metaID string
asn int64
holder string
count int64
err error
}
var valid []*store.ASEntry
for _, e := range list {
if e != nil && store.ValidASN(e.ASN) {
valid = append(valid, e)
}
}
sem := make(chan struct{}, collectConcurrency())
results := make([]entryResult, len(valid))
var wg sync.WaitGroup
for i, e := range valid {
wg.Add(1)
go func(idx int, entry *store.ASEntry) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
comm := entry.CommunityID
if comm == nil && mod.DefaultCommunityID != nil {
c := *mod.DefaultCommunityID
comm = &c
}
pfxs, holder, err := resolveASNForEntry(ctx, st, hc, entry.ASN)
if err != nil {
results[idx] = entryResult{err: fmt.Errorf("resolve AS%d: %w", entry.ASN, err)}
return
}
src := fmt.Sprintf("as:%d", entry.ASN)
var rows []store.PrefixRow
for _, pfx := range pfxs {
rows = append(rows, store.PrefixRow{Prefix: pfx.String(), CommunityID: comm, Source: src})
}
results[idx] = entryResult{
rows: rows,
metaID: entry.ID,
asn: entry.ASN,
holder: holder,
count: int64(len(pfxs)),
}
}(i, e)
}
wg.Wait()
seenPfx := make(map[string]struct{})
var out []store.PrefixRow
now := time.Now().UTC()
for _, r := range results {
if r.err != nil {
return nil, r.err
}
if r.metaID != "" {
if err := st.UpdateASEntryResolveMeta(tenantID, moduleID, r.metaID, r.holder, r.count, now); err != nil {
return nil, fmt.Errorf("as entry meta AS%d: %w", r.asn, err)
}
}
for _, row := range r.rows {
k := row.Prefix
if _, ok := seenPfx[k]; ok {
continue
}
seenPfx[k] = struct{}{}
out = append(out, row)
}
}
return out, nil
}
func collectCDNPrefixRows(ctx context.Context, st store.Backend, hc *http.Client, tenantID string, mod *store.Module, sources []*store.CDNSource, priorSnapshot []store.PrefixRow) ([]store.PrefixRow, error) {
moduleID := mod.ID
now := time.Now().UTC()
var valid []*store.CDNSource
for _, s := range sources {
if s != nil {
valid = append(valid, s)
}
}
type srcResult struct {
rows []store.PrefixRow
err error
}
results := make([]srcResult, len(valid))
sem := make(chan struct{}, collectConcurrency())
var wg sync.WaitGroup
for i, src := range valid {
wg.Add(1)
go func(idx int, src *store.CDNSource) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
sourceKey := "cdn:" + src.ID
if shouldSkipCDNSourceFetch(src, now) {
if cached := prefixRowsForSource(priorSnapshot, sourceKey); len(cached) > 0 {
results[idx] = srcResult{rows: cached}
return
}
if cached := latestCDNRowsBySource(st, tenantID)[sourceKey]; len(cached) > 0 {
results[idx] = srcResult{rows: cached}
return
}
}
rows, err := applyCDNSourceHTTPResult(ctx, st, hc, tenantID, moduleID, mod, src, priorSnapshot, now)
if err != nil {
results[idx] = srcResult{err: err}
return
}
results[idx] = srcResult{rows: rows}
}(i, src)
}
wg.Wait()
var out []store.PrefixRow
for _, r := range results {
if r.err != nil {
return nil, r.err
}
out = append(out, r.rows...)
}
return out, nil
}
func collectDomainPrefixRows(ctx context.Context, hc *http.Client, mod *store.Module, profile *store.DohProfile, entries []*store.DomainEntry) ([]store.PrefixRow, error) {
var validDom []*store.DomainEntry
for _, e := range entries {
if e != nil {
validDom = append(validDom, e)
}
}
type domResult struct {
rows []store.PrefixRow
err error
}
results := make([]domResult, len(validDom))
sem := make(chan struct{}, collectConcurrency())
var wg sync.WaitGroup
for i, e := range validDom {
wg.Add(1)
go func(idx int, entry *store.DomainEntry) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
comm := entry.CommunityID
if comm == nil && mod.DefaultCommunityID != nil {
c := *mod.DefaultCommunityID
comm = &c
}
addrs, err := resolveDomainIPs(ctx, hc, profile, entry.FQDN)
if err != nil {
results[idx] = domResult{err: fmt.Errorf("resolve domain %q: %w", entry.FQDN, err)}
return
}
src := "domain:" + strings.TrimSpace(entry.FQDN)
var rows []store.PrefixRow
for _, ip := range addrs {
cidr := ipToHostPrefix(ip)
if cidr == "" {
continue
}
rows = append(rows, store.PrefixRow{
Prefix: cidr,
CommunityID: comm,
Source: src,
})
}
results[idx] = domResult{rows: rows}
}(i, e)
}
wg.Wait()
seen := make(map[string]struct{})
var out []store.PrefixRow
for _, r := range results {
if r.err != nil {
return nil, r.err
}
for _, row := range r.rows {
key := row.Prefix + "|" + row.Source
if _, ok := seen[key]; ok {
continue
}
seen[key] = struct{}{}
out = append(out, row)
}
}
return out, nil
}