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
+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)