Files
DenozordecandCursor 8fe74c1d3b feat(jobs): add durable PG queue reclaim, slog, and richer metrics
JSON slog в ключевых пакетах; Prometheus path_group, job_audit_depth, upstream breaker; job_audit ClaimQueued/ReclaimStaleRunning + Adopt loop для HA после рестарта.

Co-authored-by: Cursor <[email protected]>
2026-07-31 12:26:18 +07:00

94 lines
2.9 KiB
Go

package pipeline
import (
"evobgp/internal/logging"
"fmt"
"net/netip"
"os"
"strings"
"evobgp/internal/store"
)
// staleOnUpstreamError reports whether ingest should keep last-known prefixes when an upstream fetch fails.
// Enabled by default; set EVOBGP_STALE_ON_UPSTREAM_ERROR=0 to restore fail-fast behavior.
func staleOnUpstreamError() bool {
v := strings.TrimSpace(os.Getenv("EVOBGP_STALE_ON_UPSTREAM_ERROR"))
if v == "" || v == "1" || strings.EqualFold(v, "true") {
return true
}
return false
}
// cdnPartialOK reports whether a failed CDN source without stale cache should be skipped
// instead of failing the whole module refresh. Opt-in: EVOBGP_CDN_PARTIAL_OK=1.
func cdnPartialOK() bool {
v := strings.TrimSpace(os.Getenv("EVOBGP_CDN_PARTIAL_OK"))
return v == "1" || strings.EqualFold(v, "true")
}
func logStaleUpstream(kind, detail string) {
logging.Default().Info(fmt.Sprintf("pipeline: stale upstream fallback (%s): %s", kind, detail))
}
func staleASNPrefixes(st store.Backend, priorSnapshot []store.PrefixRow, asn int64) ([]store.PrefixRow, string, bool) {
sourceKey := fmt.Sprintf("as:%d", asn)
if cached := prefixRowsForSource(priorSnapshot, sourceKey); len(cached) > 0 {
return cached, "", true
}
if st == nil {
return nil, "", false
}
ent, ok, err := st.GetASNPrefixCache(asn)
if err != nil || !ok || ent == nil || len(ent.Prefixes) == 0 {
return nil, "", false
}
var rows []store.PrefixRow
for _, p := range ent.Prefixes {
pfx, perr := netip.ParsePrefix(strings.TrimSpace(p))
if perr != nil {
continue
}
rows = append(rows, store.PrefixRow{Prefix: pfx.Masked().String(), Source: sourceKey})
}
if len(rows) == 0 {
return nil, "", false
}
return rows, ent.Holder, true
}
func staleDomainPrefixes(priorSnapshot []store.PrefixRow, fqdn string) ([]store.PrefixRow, bool) {
sourceKey := "domain:" + strings.TrimSpace(fqdn)
cached := prefixRowsForSource(priorSnapshot, sourceKey)
return cached, len(cached) > 0
}
func staleCDNPrefixes(st store.Backend, tenantID, moduleID string, priorSnapshot []store.PrefixRow, sourceID string) ([]store.PrefixRow, bool) {
sourceKey := cdnSourceKey(sourceID)
cached := cachedCDNPrefixRows(st, tenantID, moduleID, priorSnapshot, sourceKey)
return cached, len(cached) > 0
}
// asnCacheExpired returns cached ASN prefixes even past TTL (for stale fallback only).
func asnCacheExpired(st store.Backend, asn int64) ([]netip.Prefix, string, bool) {
if st == nil {
return nil, "", false
}
ent, ok, err := st.GetASNPrefixCache(asn)
if err != nil || !ok || ent == nil || len(ent.Prefixes) == 0 {
return nil, "", false
}
out := make([]netip.Prefix, 0, len(ent.Prefixes))
for _, p := range ent.Prefixes {
pfx, perr := netip.ParsePrefix(strings.TrimSpace(p))
if perr != nil {
continue
}
out = append(out, pfx.Masked())
}
if len(out) == 0 {
return nil, "", false
}
return out, ent.Holder, true
}