feat: enhance worker functionality by adding job queuing for deploy_apply. Introduce a Registry field in the Worker struct and implement enqueueDeployAllSpeakers method to streamline job processing after refresh or rollback operations, improving deployment efficiency.
CI / changes (push) Successful in 5s
CI / openapi (push) Has been skipped
CI / go (push) Successful in 25s
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) Failing after 12s
CI / docker-go-prime (push) Successful in 24s
CI / docker-go (deploy/docker/evobgp-agent/Dockerfile, , evobgp-agent) (push) Successful in 1m3s
CI / docker-go (evobgp-all, 1, deploy/docker/gobinary/Dockerfile, , evobgp-all) (push) Successful in 2m17s
CI / docker-go (evobgp-api, 1, deploy/docker/gobinary/Dockerfile, , evobgp-api) (push) Successful in 1m23s
CI / docker-go (evobgp-deploy, 0, deploy/docker/gobinary/Dockerfile, , evobgp-deploy) (push) Successful in 1m24s
CI / docker-go (evobgp-ingest, 0, deploy/docker/gobinary/Dockerfile, , evobgp-ingest) (push) Successful in 1m20s
CI / docker-go (evobgp-node, 0, deploy/docker/gobinary/Dockerfile, , evobgp-node) (push) Successful in 1m5s
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 1m20s

This commit is contained in:
Denozordec
2026-04-06 16:32:23 +07:00
parent e20c9f3113
commit 42d956d109
2 changed files with 26 additions and 0 deletions
+1
View File
@@ -47,6 +47,7 @@ func BootstrapWorkers(ctx context.Context, opts Options) (store.Backend, *jobs.R
cdnHTTP := &http.Client{Timeout: 45 * time.Second}
wk := &jobs.Worker{Store: backend, HTTPClient: cdnHTTP}
reg := jobs.NewRegistry(wk.Process)
wk.Registry = reg
observability.RegisterStoreBackend(backend)
return backend, reg, pool, nil
}
+25
View File
@@ -47,6 +47,8 @@ const (
type Worker struct {
Store store.Backend
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
}
var defaultWorkerHTTP = &http.Client{Timeout: 45 * time.Second}
@@ -88,6 +90,7 @@ func (w *Worker) Process(j *Job) {
return
}
j.mergeMeta(map[string]any{"revision_id": rev})
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
j.Succeed()
case KindDeployApply:
w.runDeployApply(j)
@@ -114,6 +117,27 @@ func (w *Worker) Process(j *Job) {
}
}
// 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 {
return
}
revID = strings.TrimSpace(revID)
if revID == "" {
return
}
applyJob, _, err := w.Registry.Enqueue(tenantID, KindDeployApply, nil, nil, map[string]any{
"revision_id": revID,
})
if err != nil {
j.mergeMeta(map[string]any{"deploy_apply_enqueue_error": err.Error()})
return
}
if applyJob != nil {
j.mergeMeta(map[string]any{"deploy_apply_job_id": applyJob.ID})
}
}
func (w *Worker) runDeployApply(j *Job) {
revID, _ := j.Meta["revision_id"].(string)
spk, hasSpeaker := j.Meta["speaker_id"].(string)
@@ -186,5 +210,6 @@ func (w *Worker) runRollback(j *Job) {
return
}
j.mergeMeta(map[string]any{"new_revision_id": newID})
w.enqueueDeployAllSpeakers(j, j.TenantID, newID)
j.Succeed()
}