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}
|
cdnHTTP := &http.Client{Timeout: 45 * time.Second}
|
||||||
wk := &jobs.Worker{Store: backend, HTTPClient: cdnHTTP}
|
wk := &jobs.Worker{Store: backend, HTTPClient: cdnHTTP}
|
||||||
reg := jobs.NewRegistry(wk.Process)
|
reg := jobs.NewRegistry(wk.Process)
|
||||||
|
wk.Registry = reg
|
||||||
observability.RegisterStoreBackend(backend)
|
observability.RegisterStoreBackend(backend)
|
||||||
return backend, reg, pool, nil
|
return backend, reg, pool, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -47,6 +47,8 @@ const (
|
|||||||
type Worker struct {
|
type Worker struct {
|
||||||
Store store.Backend
|
Store store.Backend
|
||||||
HTTPClient *http.Client // optional; CDN refresh uses this (default 45s timeout).
|
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}
|
var defaultWorkerHTTP = &http.Client{Timeout: 45 * time.Second}
|
||||||
@@ -88,6 +90,7 @@ func (w *Worker) Process(j *Job) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
j.mergeMeta(map[string]any{"revision_id": rev})
|
j.mergeMeta(map[string]any{"revision_id": rev})
|
||||||
|
w.enqueueDeployAllSpeakers(j, j.TenantID, rev)
|
||||||
j.Succeed()
|
j.Succeed()
|
||||||
case KindDeployApply:
|
case KindDeployApply:
|
||||||
w.runDeployApply(j)
|
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) {
|
func (w *Worker) runDeployApply(j *Job) {
|
||||||
revID, _ := j.Meta["revision_id"].(string)
|
revID, _ := j.Meta["revision_id"].(string)
|
||||||
spk, hasSpeaker := j.Meta["speaker_id"].(string)
|
spk, hasSpeaker := j.Meta["speaker_id"].(string)
|
||||||
@@ -186,5 +210,6 @@ func (w *Worker) runRollback(j *Job) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
j.mergeMeta(map[string]any{"new_revision_id": newID})
|
j.mergeMeta(map[string]any{"new_revision_id": newID})
|
||||||
|
w.enqueueDeployAllSpeakers(j, j.TenantID, newID)
|
||||||
j.Succeed()
|
j.Succeed()
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user