refactor: streamline module refresh process and enhance revision handling
CI / changes (push) Successful in 7s
CI / openapi (push) Has been skipped
CI / go (push) Successful in 40s
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 23s
CI / docker-go (deploy/docker/evobgp-agent/Dockerfile, , evobgp-agent) (push) Successful in 1m1s
CI / docker-go (evobgp-all, 1, deploy/docker/gobinary/Dockerfile, , evobgp-all) (push) Successful in 2m11s
CI / docker-go (evobgp-api, 1, deploy/docker/gobinary/Dockerfile, , evobgp-api) (push) Successful in 1m26s
CI / docker-go (evobgp-deploy, 0, deploy/docker/gobinary/Dockerfile, , evobgp-deploy) (push) Successful in 1m25s
CI / docker-go (evobgp-ingest, 0, deploy/docker/gobinary/Dockerfile, , evobgp-ingest) (push) Successful in 1m24s
CI / docker-go (evobgp-node, 0, deploy/docker/gobinary/Dockerfile, , evobgp-node) (push) Successful in 1m9s
CI / docker-go (evobgp-render, 0, deploy/docker/gobinary/Dockerfile, , evobgp-render) (push) Successful in 1m21s
CI / docker-go (evobgp-scheduler, 0, deploy/docker/gobinary/Dockerfile, , evobgp-scheduler) (push) Successful in 1m21s
CI / changes (push) Successful in 7s
CI / openapi (push) Has been skipped
CI / go (push) Successful in 40s
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 23s
CI / docker-go (deploy/docker/evobgp-agent/Dockerfile, , evobgp-agent) (push) Successful in 1m1s
CI / docker-go (evobgp-all, 1, deploy/docker/gobinary/Dockerfile, , evobgp-all) (push) Successful in 2m11s
CI / docker-go (evobgp-api, 1, deploy/docker/gobinary/Dockerfile, , evobgp-api) (push) Successful in 1m26s
CI / docker-go (evobgp-deploy, 0, deploy/docker/gobinary/Dockerfile, , evobgp-deploy) (push) Successful in 1m25s
CI / docker-go (evobgp-ingest, 0, deploy/docker/gobinary/Dockerfile, , evobgp-ingest) (push) Successful in 1m24s
CI / docker-go (evobgp-node, 0, deploy/docker/gobinary/Dockerfile, , evobgp-node) (push) Successful in 1m9s
CI / docker-go (evobgp-render, 0, deploy/docker/gobinary/Dockerfile, , evobgp-render) (push) Successful in 1m21s
CI / docker-go (evobgp-scheduler, 0, deploy/docker/gobinary/Dockerfile, , evobgp-scheduler) (push) Successful in 1m21s
Updated the worker's job processing to utilize a new `RefreshModuleIngest` function, which simplifies the module refresh logic by separating ingestion from revision rendering. The `finishModuleRefreshSuccess` method now handles the aggregation of revision logs and enqueues deploy actions more efficiently. Additionally, refactored the `RefreshModule` function to maintain backward compatibility while improving clarity and functionality in the pipeline's revision management.
This commit is contained in:
+30
-22
@@ -100,22 +100,11 @@ func (w *Worker) Process(j *Job) {
|
|||||||
j.Fail("missing module_id in job meta")
|
j.Fail("missing module_id in job meta")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
rev, err := pipeline.RefreshModule(context.Background(), w.Store, w.httpClient(), j.TenantID, mid)
|
if err := pipeline.RefreshModuleIngest(context.Background(), w.Store, w.httpClient(), j.TenantID, mid); err != nil {
|
||||||
if err != nil {
|
|
||||||
j.Fail(err.Error())
|
j.Fail(err.Error())
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
j.mergeMeta(map[string]any{"revision_id": rev})
|
w.finishModuleRefreshSuccess(j, mid)
|
||||||
if entries, total, err := w.buildRevisionLogEntries(j.TenantID, rev); 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()})
|
|
||||||
}
|
|
||||||
w.finishModuleRefreshSuccess(j, rev)
|
|
||||||
case KindDeployApply:
|
case KindDeployApply:
|
||||||
w.runDeployApply(j)
|
w.runDeployApply(j)
|
||||||
case KindRevisionRollback:
|
case KindRevisionRollback:
|
||||||
@@ -146,27 +135,46 @@ func (w *Worker) tenantRefreshMu(tenantID string) *sync.Mutex {
|
|||||||
return v.(*sync.Mutex)
|
return v.(*sync.Mutex)
|
||||||
}
|
}
|
||||||
|
|
||||||
// finishModuleRefreshSuccess marks the job succeeded and enqueues deploy_apply only when no other
|
// finishModuleRefreshSuccess marks the refresh job and, for the last active refresh in tenant,
|
||||||
// module_refresh is still queued or running for the same tenant (coalesces parallel refreshes).
|
// creates one aggregate revision and enqueues a single deploy_apply.
|
||||||
func (w *Worker) finishModuleRefreshSuccess(j *Job, rev string) {
|
func (w *Worker) finishModuleRefreshSuccess(j *Job, triggerModuleID string) {
|
||||||
if w == nil || w.Registry == nil {
|
if w == nil || w.Store == nil {
|
||||||
j.Succeed()
|
j.Succeed()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
mu := w.tenantRefreshMu(j.TenantID)
|
mu := w.tenantRefreshMu(j.TenantID)
|
||||||
mu.Lock()
|
mu.Lock()
|
||||||
deferDeploy := w.Registry.CountOtherActiveModuleRefresh(j.TenantID, j.ID) > 0
|
defer mu.Unlock()
|
||||||
|
deferDeploy := false
|
||||||
|
if w.Registry != nil {
|
||||||
|
deferDeploy = w.Registry.CountOtherActiveModuleRefresh(j.TenantID, j.ID) > 0
|
||||||
|
}
|
||||||
if deferDeploy {
|
if deferDeploy {
|
||||||
j.mergeMeta(map[string]any{
|
j.mergeMeta(map[string]any{
|
||||||
"deploy_apply_deferred": true,
|
"deploy_apply_deferred": true,
|
||||||
"deploy_apply_defer_reason": "parallel_module_refresh",
|
"deploy_apply_defer_reason": "parallel_module_refresh",
|
||||||
})
|
})
|
||||||
|
j.Succeed()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
rev, err := pipeline.RenderTenantRevision(context.Background(), w.Store, w.httpClient(), j.TenantID, triggerModuleID)
|
||||||
|
if err != nil {
|
||||||
|
j.Fail(err.Error())
|
||||||
|
return
|
||||||
|
}
|
||||||
|
j.mergeMeta(map[string]any{"revision_id": rev})
|
||||||
|
if entries, total, err := w.buildRevisionLogEntries(j.TenantID, rev); 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()
|
j.Succeed()
|
||||||
mu.Unlock()
|
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
|
||||||
if !deferDeploy {
|
|
||||||
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// enqueueDeployAllSpeakers queues the same work as POST /v1/apply (all speakers, no speaker_id).
|
// enqueueDeployAllSpeakers queues the same work as POST /v1/apply (all speakers, no speaker_id).
|
||||||
|
|||||||
@@ -37,28 +37,34 @@ func MaterializedASPrefixKey(asn int64) string {
|
|||||||
return fmt.Sprintf("as:%d", asn)
|
return fmt.Sprintf("as:%d", asn)
|
||||||
}
|
}
|
||||||
|
|
||||||
// RefreshModule runs ingest (where applicable) for one module, then renders a new revision whose
|
// RefreshModuleIngest runs ingest for one module and persists side-effects (ASN metadata, CDN etags, etc).
|
||||||
// BIRD materialization includes prefixes from all enabled modules of the tenant (others via live collect).
|
// It does not create a new config revision.
|
||||||
// If the tenant-wide materialized prefix set is unchanged from the latest revision, returns that
|
func RefreshModuleIngest(ctx context.Context, st store.Backend, hc *http.Client, tenantID, moduleID string) error {
|
||||||
// revision id and does not insert a duplicate config_revision.
|
|
||||||
func RefreshModule(ctx context.Context, st store.Backend, hc *http.Client, tenantID, moduleID string) (revisionID string, err error) {
|
|
||||||
if hc == nil {
|
if hc == nil {
|
||||||
hc = http.DefaultClient
|
hc = http.DefaultClient
|
||||||
}
|
}
|
||||||
mod, err := st.GetModule(tenantID, moduleID)
|
mod, err := st.GetModule(tenantID, moduleID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return err
|
||||||
}
|
}
|
||||||
if !mod.Enabled {
|
if !mod.Enabled {
|
||||||
return "", fmt.Errorf("module disabled")
|
return fmt.Errorf("module disabled")
|
||||||
}
|
}
|
||||||
|
|
||||||
rows, err := collectModulePrefixRows(ctx, st, hc, tenantID, mod)
|
_, err = collectModulePrefixRows(ctx, st, hc, tenantID, mod)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return err
|
||||||
}
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
agg, err := aggregateTenantPrefixRows(ctx, st, hc, tenantID, moduleID, rows)
|
// RenderTenantRevision renders one tenant-wide revision using current data from all enabled modules.
|
||||||
|
// If materialized prefixes are unchanged, returns latest revision id without creating a duplicate.
|
||||||
|
func RenderTenantRevision(ctx context.Context, st store.Backend, hc *http.Client, tenantID, triggerModuleID string) (revisionID string, err error) {
|
||||||
|
if hc == nil {
|
||||||
|
hc = http.DefaultClient
|
||||||
|
}
|
||||||
|
agg, err := aggregateTenantPrefixRowsAll(ctx, st, hc, tenantID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
@@ -69,17 +75,25 @@ func RefreshModule(ctx context.Context, st store.Backend, hc *http.Client, tenan
|
|||||||
agg = smartAggregatePrefixRows(agg)
|
agg = smartAggregatePrefixRows(agg)
|
||||||
|
|
||||||
revisionID = uuid.NewString()
|
revisionID = uuid.NewString()
|
||||||
parent := parentRevision(st, tenantID, moduleID)
|
parent := parentRevision(st, tenantID, triggerModuleID)
|
||||||
preview, err := buildPreviewFragments(st, tenantID, moduleID, revisionID, agg)
|
preview, err := buildPreviewFragments(st, tenantID, triggerModuleID, revisionID, agg)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
if err := st.CreateRenderRevision(revisionID, tenantID, moduleID, parent, hash, preview, agg); err != nil {
|
if err := st.CreateRenderRevision(revisionID, tenantID, triggerModuleID, parent, hash, preview, agg); err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
return revisionID, nil
|
return revisionID, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// RefreshModule keeps backwards-compatible behavior: module ingest + immediate tenant render.
|
||||||
|
func RefreshModule(ctx context.Context, st store.Backend, hc *http.Client, tenantID, moduleID string) (revisionID string, err error) {
|
||||||
|
if err := RefreshModuleIngest(ctx, st, hc, tenantID, moduleID); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return RenderTenantRevision(ctx, st, hc, tenantID, moduleID)
|
||||||
|
}
|
||||||
|
|
||||||
// collectModulePrefixRows returns materialized prefix rows for a single module (source of truth from store / ASN resolve / CDN fetch).
|
// collectModulePrefixRows returns materialized prefix rows for a single module (source of truth from store / ASN resolve / CDN fetch).
|
||||||
func collectModulePrefixRows(ctx context.Context, st store.Backend, hc *http.Client, tenantID string, mod *store.Module) ([]store.PrefixRow, error) {
|
func collectModulePrefixRows(ctx context.Context, st store.Backend, hc *http.Client, tenantID string, mod *store.Module) ([]store.PrefixRow, error) {
|
||||||
moduleID := mod.ID
|
moduleID := mod.ID
|
||||||
@@ -452,21 +466,15 @@ func ipToHostPrefix(ip netip.Addr) string {
|
|||||||
return netip.PrefixFrom(ip, bits).Masked().String()
|
return netip.PrefixFrom(ip, bits).Masked().String()
|
||||||
}
|
}
|
||||||
|
|
||||||
// aggregateTenantPrefixRows builds the union of materialized prefixes for all enabled modules.
|
// aggregateTenantPrefixRowsAll builds the union of materialized prefixes for all enabled modules
|
||||||
// The module that triggered refresh contributes freshRows; every other module is collected live from the store
|
// using current source data from store/external resolvers.
|
||||||
// (same logic as refresh). We do not reuse other modules' saved revisions as prefix sources, because each revision
|
func aggregateTenantPrefixRowsAll(ctx context.Context, st store.Backend, hc *http.Client, tenantID string) ([]store.PrefixRow, error) {
|
||||||
// already stores the full tenant-wide aggregate — mixing them with freshRows would duplicate prefixes.
|
|
||||||
func aggregateTenantPrefixRows(ctx context.Context, st store.Backend, hc *http.Client, tenantID, changedModuleID string, freshRows []store.PrefixRow) ([]store.PrefixRow, error) {
|
|
||||||
mods := st.ListModules(tenantID)
|
mods := st.ListModules(tenantID)
|
||||||
var out []store.PrefixRow
|
var out []store.PrefixRow
|
||||||
for _, m := range mods {
|
for _, m := range mods {
|
||||||
if m == nil || !m.Enabled {
|
if m == nil || !m.Enabled {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if m.ID == changedModuleID {
|
|
||||||
out = append(out, freshRows...)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
omod, err := st.GetModule(tenantID, m.ID)
|
omod, err := st.GetModule(tenantID, m.ID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|||||||
Reference in New Issue
Block a user