From 42d956d109e3a9634db65aa18dd9f676da80b3d8 Mon Sep 17 00:00:00 2001 From: Denozordec Date: Mon, 6 Apr 2026 16:32:23 +0700 Subject: [PATCH] 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. --- internal/httpapi/bootstrap.go | 1 + internal/jobs/worker.go | 25 +++++++++++++++++++++++++ 2 files changed, 26 insertions(+) diff --git a/internal/httpapi/bootstrap.go b/internal/httpapi/bootstrap.go index 566b2d4..256f621 100644 --- a/internal/httpapi/bootstrap.go +++ b/internal/httpapi/bootstrap.go @@ -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 } diff --git a/internal/jobs/worker.go b/internal/jobs/worker.go index 2f900e3..5eb910a 100644 --- a/internal/jobs/worker.go +++ b/internal/jobs/worker.go @@ -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() }