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

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:
Denozordec
2026-04-08 13:42:53 +07:00
parent 2260ecf74a
commit 10e2415cfa
2 changed files with 60 additions and 44 deletions
+30 -22
View File
@@ -100,22 +100,11 @@ func (w *Worker) Process(j *Job) {
j.Fail("missing module_id in job meta")
return
}
rev, err := pipeline.RefreshModule(context.Background(), w.Store, w.httpClient(), j.TenantID, mid)
if err != nil {
if err := pipeline.RefreshModuleIngest(context.Background(), w.Store, w.httpClient(), j.TenantID, mid); 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()})
}
w.finishModuleRefreshSuccess(j, rev)
w.finishModuleRefreshSuccess(j, mid)
case KindDeployApply:
w.runDeployApply(j)
case KindRevisionRollback:
@@ -146,27 +135,46 @@ func (w *Worker) tenantRefreshMu(tenantID string) *sync.Mutex {
return v.(*sync.Mutex)
}
// finishModuleRefreshSuccess marks the job succeeded and enqueues deploy_apply only when no other
// module_refresh is still queued or running for the same tenant (coalesces parallel refreshes).
func (w *Worker) finishModuleRefreshSuccess(j *Job, rev string) {
if w == nil || w.Registry == nil {
// finishModuleRefreshSuccess marks the refresh job and, for the last active refresh in tenant,
// creates one aggregate revision and enqueues a single deploy_apply.
func (w *Worker) finishModuleRefreshSuccess(j *Job, triggerModuleID string) {
if w == nil || w.Store == nil {
j.Succeed()
return
}
mu := w.tenantRefreshMu(j.TenantID)
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 {
j.mergeMeta(map[string]any{
"deploy_apply_deferred": true,
"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()
mu.Unlock()
if !deferDeploy {
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
}
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
}
// enqueueDeployAllSpeakers queues the same work as POST /v1/apply (all speakers, no speaker_id).