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
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:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user