feat: implement peer reconciliation job for BGP peers
CI / changes (push) Successful in 7s
CI / openapi (push) Successful in 23s
CI / go (push) Failing after 30s
CI / docker-web (deploy/docker/evobgp-web/Dockerfile, , evobgp-web) (push) Successful in 1m5s
CI / docker-web (deploy/docker/evobgp-web/Dockerfile, evobgp-all, evobgp-web-all) (push) Successful in 1m3s
CI / docker-bird (push) Has been skipped
CI / bird2 (push) Has been skipped
CI / docker-go-prime (push) Has been skipped
CI / docker-go (deploy/docker/evobgp-agent/Dockerfile, , evobgp-agent) (push) Has been skipped
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

Added a new job type `peer_reconcile` to handle fast reconciliation of BGP peers without module ingestion. Updated API documentation to reflect the new job's functionality and its automatic application to speakers post-reconciliation. Enhanced the HTTP API to enqueue reconciliation jobs during peer creation, updates, and deletions, ensuring efficient state management for BGP configurations.
This commit is contained in:
Denozordec
2026-04-09 15:07:21 +07:00
parent e700f90c47
commit 3b17228ef2
8 changed files with 289 additions and 5 deletions
+1
View File
@@ -48,6 +48,7 @@
- `GET /v1/peers`, `POST /v1/peers`
- `GET|PATCH|DELETE /v1/peers/{id}`
- Для `POST|PATCH|DELETE` peer запускается быстрый job `peer_reconcile` (без module ingest/сбора префиксов); после него автоматически ставится apply на спикеры.
### Speakers
+11 -1
View File
@@ -185,7 +185,9 @@ components:
in: query
schema:
type: string
description: Фильтр по виду задачи; точный перечень расширяем.
description: >
Фильтр по виду задачи; точный перечень расширяем.
Основные значения: `module_refresh`, `peer_reconcile`, `deploy_apply`, `revision_rollback`, `bird_reload`.
responses:
Unauthorized:
@@ -1957,6 +1959,9 @@ paths:
post:
tags: [Peers]
summary: Создать пира
description: >
Создаёт BGP-пира и инициирует быстрый reconcile пиров (job `peer_reconcile`) без module ingest.
После reconcile автоматически запускается apply на спикеры.
operationId: createPeer
parameters:
- $ref: "#/components/parameters/TenantId"
@@ -2002,6 +2007,8 @@ paths:
description: >
Политики (`policies_json`: `local_ipv4`, `local_ipv6`, `local_asn`), neighbor, ASN,
привязка к `bgp_speaker_id` или `null` для всех спикеров.
Изменение инициирует быстрый reconcile пиров (job `peer_reconcile`) без module ingest
и затем авто-apply на спикеры.
operationId: patchPeer
parameters:
- $ref: "#/components/parameters/IdempotencyKey"
@@ -2025,6 +2032,9 @@ paths:
delete:
tags: [Peers]
summary: Удалить или отключить пира
description: >
Удаление/отключение инициирует быстрый reconcile пиров (job `peer_reconcile`) без module ingest
и затем авто-apply на спикеры.
operationId: deletePeer
parameters:
- $ref: "#/components/parameters/IdempotencyKey"
+11
View File
@@ -451,6 +451,17 @@ func (s *Server) enqueueModuleRefreshIfEnabled(tenantID, moduleID, trigger strin
})
}
// enqueuePeerReconcile queues fast peer-only reconcile/render (best-effort, no HTTP error).
func (s *Server) enqueuePeerReconcile(tenantID, trigger string) {
if s.jobs == nil {
return
}
_, _, _ = s.jobs.Enqueue(tenantID, jobs.KindPeerReconcile, nil, nil, map[string]any{
"trigger": trigger,
"job_title": "Обновление BGP пиров",
})
}
func (s *Server) handleGetRevision(w http.ResponseWriter, r *http.Request) {
a, ok := authFromContext(r.Context())
if !ok {
+3
View File
@@ -1028,6 +1028,7 @@ func (s *Server) handlePostPeer(w http.ResponseWriter, r *http.Request) {
writeStoreErr(w, err)
return
}
s.enqueuePeerReconcile(a.TenantID, "peer_create")
writeJSON(w, http.StatusCreated, peerJSON(x))
}
@@ -1059,6 +1060,7 @@ func (s *Server) handlePatchPeer(w http.ResponseWriter, r *http.Request) {
writeStoreErr(w, err)
return
}
s.enqueuePeerReconcile(a.TenantID, "peer_patch")
writeJSON(w, http.StatusOK, peerJSON(x))
}
@@ -1071,6 +1073,7 @@ func (s *Server) handleDeletePeer(w http.ResponseWriter, r *http.Request) {
writeStoreErr(w, err)
return
}
s.enqueuePeerReconcile(a.TenantID, "peer_delete")
w.WriteHeader(http.StatusNoContent)
}
+84
View File
@@ -42,6 +42,7 @@ func mergeBirdPostApplyMeta(j *Job) {
const (
KindModuleRefresh = "module_refresh"
KindPeerReconcile = "peer_reconcile"
KindDeployApply = "deploy_apply"
KindRevisionRollback = "revision_rollback"
KindBirdReload = "bird_reload"
@@ -105,6 +106,8 @@ func (w *Worker) Process(j *Job) {
return
}
w.finishModuleRefreshSuccess(j, mid)
case KindPeerReconcile:
w.runPeerReconcile(j)
case KindDeployApply:
w.runDeployApply(j)
case KindRevisionRollback:
@@ -130,6 +133,87 @@ func (w *Worker) Process(j *Job) {
}
}
func (w *Worker) runPeerReconcile(j *Job) {
if w == nil || w.Store == nil {
j.Fail("worker not configured")
return
}
const peerJobTitle = "Обновление BGP пиров"
j.mergeMeta(map[string]any{"job_title": peerJobTitle})
var revID string
latest, _, _ := w.Store.ListRevisions(j.TenantID, "", "", 1)
triggerModuleID, err := w.peerTriggerModuleID(j.TenantID, latest)
if err != nil {
j.Fail(err.Error())
return
}
if len(latest) == 0 {
// First run fallback: render full tenant state once if no baseline revision exists yet.
rid, err := pipeline.RenderTenantRevision(context.Background(), w.Store, w.httpClient(), j.TenantID, triggerModuleID)
if err != nil {
j.Fail(err.Error())
return
}
revID = rid
} else {
baseRevID := latest[0].ID
rows := make([]store.PrefixRow, 0, 1024)
cursor := ""
for {
page, next, more := w.Store.ListRevisionPrefixes(j.TenantID, baseRevID, cursor, 2000)
rows = append(rows, page...)
if !more || strings.TrimSpace(next) == "" {
break
}
cursor = next
}
rid, err := pipeline.RenderTenantRevisionFromPrefixes(context.Background(), w.Store, w.httpClient(), j.TenantID, triggerModuleID, rows)
if err != nil {
j.Fail(err.Error())
return
}
revID = rid
}
j.mergeMeta(map[string]any{"revision_id": revID})
if entries, total, err := w.buildRevisionLogEntries(j.TenantID, revID); err == nil {
j.mergeMeta(map[string]any{
"log_entries": entries,
"log_total": total,
"log_generated": time.Now().UTC().Format(time.RFC3339Nano),
})
} else {
j.mergeMeta(map[string]any{"log_build_error": err.Error()})
}
j.Succeed()
w.enqueueDeployAllSpeakers(j, j.TenantID, revID)
}
func (w *Worker) peerTriggerModuleID(tenantID string, latest []*store.Revision) (string, error) {
if len(latest) > 0 {
if mid := strings.TrimSpace(latest[0].ModuleID); mid != "" {
return mid, nil
}
}
for _, mod := range w.Store.ListModules(tenantID) {
if mod == nil || !mod.Enabled {
continue
}
if strings.TrimSpace(mod.ID) != "" {
return mod.ID, nil
}
}
for _, mod := range w.Store.ListModules(tenantID) {
if mod == nil {
continue
}
if strings.TrimSpace(mod.ID) != "" {
return mod.ID, nil
}
}
return "", fmt.Errorf("missing module_id for peer reconcile")
}
func (w *Worker) tenantRefreshMu(tenantID string) *sync.Mutex {
v, _ := w.refreshGate.LoadOrStore(tenantID, &sync.Mutex{})
return v.(*sync.Mutex)
+86 -3
View File
@@ -48,7 +48,7 @@ func TestParallelModuleRefresh_CoalescesDeployApply(t *testing.T) {
t.Fatal(err)
}
waitSucceededModuleRefreshCount(t, reg, tenant, 2)
waitSucceededJobsByKindCount(t, reg, tenant, KindModuleRefresh, 2)
deployJobs, _, _ := reg.List(tenant, "", KindDeployApply, "", 100)
if len(deployJobs) != 1 {
@@ -75,14 +75,97 @@ func TestParallelModuleRefresh_CoalescesDeployApply(t *testing.T) {
}
func waitSucceededModuleRefreshCount(t *testing.T, reg *Registry, tenant string, want int) {
t.Helper()
waitSucceededJobsByKindCount(t, reg, tenant, KindModuleRefresh, want)
}
func waitSucceededJobsByKindCount(t *testing.T, reg *Registry, tenant, kind string, want int) {
t.Helper()
deadline := time.Now().Add(30 * time.Second)
for time.Now().Before(deadline) {
jobs, _, _ := reg.List(tenant, StatusSucceeded, KindModuleRefresh, "", 100)
jobs, _, _ := reg.List(tenant, StatusSucceeded, kind, "", 100)
if len(jobs) >= want {
return
}
all, _, _ := reg.List(tenant, "", kind, "", 100)
for _, j := range all {
if j == nil || j.Status != StatusFailed {
continue
}
snap := j.Snapshot()
t.Fatalf("%s job failed: %#v", kind, snap["error"])
}
time.Sleep(5 * time.Millisecond)
}
t.Fatal("timeout waiting for module_refresh jobs")
t.Fatalf("timeout waiting for %s jobs", kind)
}
func TestPeerReconcile_RendersFromLatestRevisionAndQueuesDeploy(t *testing.T) {
t.Setenv("EVOBGP_ASN_RESOLVE", "0")
t.Setenv("EVOBGP_BIRD_ACTIVE_DIR", "") // skip bird binary path in deploy_apply
m := store.NewMemory()
m.SeedDemo()
tenant, _, modIP, _, _ := m.DemoIDs()
for _, mod := range m.ListModules(tenant) {
if mod == nil || mod.ID == modIP {
continue
}
disabled := false
if _, err := m.UpdateModule(tenant, mod.ID, &store.ModulePatch{Enabled: &disabled}); err != nil {
t.Fatal(err)
}
}
if _, err := m.CreateIPRangeEntry(tenant, modIP, &store.IPRangeEntry{Prefix: "10.10.0.0/24"}); err != nil {
t.Fatal(err)
}
peer, err := m.CreatePeer(tenant, &store.BGPPeer{
Name: "peer-a",
Neighbor: "192.0.2.2",
RemoteASN: 64512,
Enabled: true,
})
if err != nil {
t.Fatal(err)
}
wk := &Worker{Store: m}
reg := NewRegistry(wk.Process)
wk.Registry = reg
// Create baseline revision from enabled IP module.
mid := modIP
if _, _, err := reg.Enqueue(tenant, KindModuleRefresh, nil, &mid, map[string]any{"module_id": modIP}); err != nil {
t.Fatal(err)
}
waitSucceededJobsByKindCount(t, reg, tenant, KindModuleRefresh, 1)
beforeRevs, _, _ := m.ListRevisions(tenant, "", "", 200)
enabled := false
if _, err := m.UpdatePeer(tenant, peer.ID, &store.PeerPatch{Enabled: &enabled}); err != nil {
t.Fatal(err)
}
if _, _, err := reg.Enqueue(tenant, KindPeerReconcile, nil, nil, map[string]any{"trigger": "peer_patch"}); err != nil {
t.Fatal(err)
}
waitSucceededJobsByKindCount(t, reg, tenant, KindPeerReconcile, 1)
afterRevs, _, _ := m.ListRevisions(tenant, "", "", 200)
if len(afterRevs) <= len(beforeRevs) {
t.Fatalf("expected new revision after peer_reconcile, before=%d after=%d", len(beforeRevs), len(afterRevs))
}
peerJobs, _, _ := reg.List(tenant, StatusSucceeded, KindPeerReconcile, "", 10)
if len(peerJobs) == 0 {
t.Fatal("expected succeeded peer_reconcile job")
}
meta := peerJobs[0].Snapshot()["meta"].(map[string]any)
if _, ok := meta["revision_id"].(string); !ok {
t.Fatalf("expected revision_id in peer_reconcile meta, got %#v", meta)
}
if title, ok := meta["job_title"].(string); !ok || title != "Обновление BGP пиров" {
t.Fatalf("expected job_title in peer_reconcile meta, got %#v", meta["job_title"])
}
if _, ok := meta["deploy_apply_job_id"].(string); !ok {
t.Fatalf("expected deploy_apply_job_id in peer_reconcile meta, got %#v", meta)
}
}
+91 -1
View File
@@ -72,7 +72,34 @@ func RenderTenantRevision(ctx context.Context, st store.Backend, hc *http.Client
if err != nil {
return "", err
}
hash := hashAggregatedMaterialization(tenantID, agg)
hash := hashAggregatedMaterializationWithPeers(st, tenantID, agg)
if prev := latestTenantRevision(st, tenantID); prev != nil && prev.ContentHash == hash {
return prev.ID, nil
}
agg = smartAggregatePrefixRows(agg)
revisionID = uuid.NewString()
parent := parentRevision(st, tenantID, triggerModuleID)
preview, err := buildPreviewFragments(st, tenantID, triggerModuleID, revisionID, agg)
if err != nil {
return "", err
}
if err := st.CreateRenderRevision(revisionID, tenantID, triggerModuleID, parent, hash, preview, agg); err != nil {
return "", err
}
applyRevisionRetention(st, tenantID)
return revisionID, nil
}
// RenderTenantRevisionFromPrefixes renders one tenant-wide revision from already materialized prefixes.
// This is used for fast paths (e.g. peer-only changes) to avoid ingest/external fetches.
func RenderTenantRevisionFromPrefixes(ctx context.Context, st store.Backend, hc *http.Client, tenantID, triggerModuleID string, rows []store.PrefixRow) (revisionID string, err error) {
_ = ctx
if hc == nil {
hc = http.DefaultClient
}
agg := append([]store.PrefixRow(nil), rows...)
hash := hashAggregatedMaterializationWithPeers(st, tenantID, agg)
if prev := latestTenantRevision(st, tenantID); prev != nil && prev.ContentHash == hash {
return prev.ID, nil
}
@@ -679,6 +706,69 @@ func hashAggregatedMaterialization(tenantID string, rows []store.PrefixRow) stri
return fmt.Sprintf("sha256:%x", h.Sum(nil))
}
func hashAggregatedMaterializationWithPeers(st store.Backend, tenantID string, rows []store.PrefixRow) string {
type line struct{ p, c, s string }
var lines []line
for _, r := range rows {
c := ""
if r.CommunityID != nil {
c = *r.CommunityID
}
lines = append(lines, line{r.Prefix, c, r.Source})
}
sort.Slice(lines, func(i, j int) bool {
if lines[i].p != lines[j].p {
return lines[i].p < lines[j].p
}
if lines[i].c != lines[j].c {
return lines[i].c < lines[j].c
}
return lines[i].s < lines[j].s
})
h := sha256.New()
h.Write([]byte(strings.TrimSpace(tenantID)))
h.Write([]byte{0})
for _, l := range lines {
h.Write([]byte(l.p))
h.Write([]byte{1})
h.Write([]byte(l.c))
h.Write([]byte{1})
h.Write([]byte(l.s))
h.Write([]byte{0})
}
h.Write([]byte("peers"))
h.Write([]byte{0})
peers := st.ListPeers(tenantID)
sort.Slice(peers, func(i, j int) bool {
if peers[i] == nil || peers[j] == nil {
return i < j
}
return peers[i].ID < peers[j].ID
})
for _, p := range peers {
if p == nil {
continue
}
speakerID := ""
if p.SpeakerID != nil {
speakerID = strings.TrimSpace(*p.SpeakerID)
}
h.Write([]byte(strings.TrimSpace(p.ID)))
h.Write([]byte{1})
h.Write([]byte(strings.TrimSpace(p.Neighbor)))
h.Write([]byte{1})
h.Write([]byte(strconv.FormatInt(p.RemoteASN, 10)))
h.Write([]byte{1})
h.Write([]byte(strconv.FormatBool(p.Enabled)))
h.Write([]byte{1})
h.Write([]byte(strings.TrimSpace(p.PoliciesJSON)))
h.Write([]byte{1})
h.Write([]byte(speakerID))
h.Write([]byte{0})
}
return fmt.Sprintf("sha256:%x", h.Sum(nil))
}
func buildPreviewFragments(st store.Backend, tenantID, moduleID, revisionID string, rows []store.PrefixRow) (map[string]string, error) {
v4, v6, pathASNs, staticGroups, err := materializeRowsForBird(st, tenantID, rows)
if err != nil {
+2
View File
@@ -17,6 +17,8 @@ export function jobKindTitle(job: JobRow, moduleNameById?: ReadonlyMap<string, s
}
case 'deploy_apply':
return 'Применение конфигурации';
case 'peer_reconcile':
return 'Обновление BGP пиров';
case 'revision_rollback':
return 'Откат ревизии';
case 'bird_reload':