feat: implement module refresh job management with concurrency control. Add CountOtherActiveModuleRefresh method to track active jobs and enhance finishModuleRefreshSuccess to defer deploy_apply when parallel refreshes are detected, improving job processing efficiency.
CI / changes (push) Successful in 6s
CI / openapi (push) Has been skipped
CI / go (push) Failing after 42s
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) 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
CI / changes (push) Successful in 6s
CI / openapi (push) Has been skipped
CI / go (push) Failing after 42s
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) 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
This commit is contained in:
+32
-2
@@ -8,6 +8,7 @@ import (
|
||||
"os"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"evobgp/internal/birddeploy"
|
||||
@@ -52,6 +53,8 @@ type Worker struct {
|
||||
HTTPClient *http.Client // optional; CDN refresh uses this (default 45s timeout).
|
||||
// Registry is set after BootstrapWorkers creates the job queue; used to chain deploy_apply after refresh/rollback.
|
||||
Registry *Registry
|
||||
// refreshGate serializes deploy_apply gating after module_refresh per tenant (see finishModuleRefreshSuccess).
|
||||
refreshGate sync.Map // map[string]*sync.Mutex
|
||||
}
|
||||
|
||||
type revisionLogEntry struct {
|
||||
@@ -111,8 +114,7 @@ func (w *Worker) Process(j *Job) {
|
||||
} else {
|
||||
j.mergeMeta(map[string]any{"log_build_error": err.Error()})
|
||||
}
|
||||
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
|
||||
j.Succeed()
|
||||
w.finishModuleRefreshSuccess(j, rev)
|
||||
case KindDeployApply:
|
||||
w.runDeployApply(j)
|
||||
case KindRevisionRollback:
|
||||
@@ -138,6 +140,34 @@ func (w *Worker) Process(j *Job) {
|
||||
}
|
||||
}
|
||||
|
||||
func (w *Worker) tenantRefreshMu(tenantID string) *sync.Mutex {
|
||||
v, _ := w.refreshGate.LoadOrStore(tenantID, &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 {
|
||||
j.Succeed()
|
||||
return
|
||||
}
|
||||
mu := w.tenantRefreshMu(j.TenantID)
|
||||
mu.Lock()
|
||||
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()
|
||||
mu.Unlock()
|
||||
if !deferDeploy {
|
||||
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
|
||||
}
|
||||
}
|
||||
|
||||
// enqueueDeployAllSpeakers queues the same work as POST /v1/apply (all speakers, no speaker_id).
|
||||
func (w *Worker) enqueueDeployAllSpeakers(j *Job, tenantID, revID string) {
|
||||
if w == nil || w.Registry == nil {
|
||||
|
||||
Reference in New Issue
Block a user