fix(pipeline): harden upstream resilience and production shutdown
DoH через DoWithRetry; CDN preview через UpstreamHTTPDo; частичный fail CDN (EVOBGP_CDN_PARTIAL_OK); безопасный доступ к Job.Meta; drain jobs при SIGTERM; ValidateProductionEnforce при EVOBGP_PRODUCTION=1. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
package jobs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"sort"
|
||||
@@ -111,6 +112,42 @@ func (j *Job) mergeMeta(kv map[string]any) {
|
||||
}
|
||||
}
|
||||
|
||||
// metaString returns a string meta field under lock (safe vs concurrent mergeMeta).
|
||||
func (j *Job) metaString(key string) string {
|
||||
j.mu.Lock()
|
||||
defer j.mu.Unlock()
|
||||
if j.Meta == nil {
|
||||
return ""
|
||||
}
|
||||
s, _ := j.Meta[key].(string)
|
||||
return s
|
||||
}
|
||||
|
||||
// metaBool returns a bool meta field under lock.
|
||||
func (j *Job) metaBool(key string) bool {
|
||||
j.mu.Lock()
|
||||
defer j.mu.Unlock()
|
||||
if j.Meta == nil {
|
||||
return false
|
||||
}
|
||||
b, _ := j.Meta[key].(bool)
|
||||
return b
|
||||
}
|
||||
|
||||
// metaCopy returns a shallow copy of job meta under lock.
|
||||
func (j *Job) metaCopy() map[string]any {
|
||||
j.mu.Lock()
|
||||
defer j.mu.Unlock()
|
||||
if j.Meta == nil {
|
||||
return nil
|
||||
}
|
||||
out := make(map[string]any, len(j.Meta))
|
||||
for k, v := range j.Meta {
|
||||
out[k] = v
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// statusLocked is used by the worker defer for metrics (any stable terminal or in-flight status).
|
||||
func (j *Job) statusLocked() string {
|
||||
j.mu.Lock()
|
||||
@@ -528,3 +565,65 @@ func (r *Registry) RequestCancel(tenantID, jobID string) (*Job, error) {
|
||||
j.RequestCancel()
|
||||
return j, nil
|
||||
}
|
||||
|
||||
// RequestCancelAll requests cancellation of all non-terminal jobs (all tenants).
|
||||
func (r *Registry) RequestCancelAll() int {
|
||||
if r == nil {
|
||||
return 0
|
||||
}
|
||||
r.mu.RLock()
|
||||
defer r.mu.RUnlock()
|
||||
n := 0
|
||||
for _, j := range r.byID {
|
||||
if j == nil {
|
||||
continue
|
||||
}
|
||||
st := j.statusLocked()
|
||||
if st == StatusSucceeded || st == StatusFailed || st == StatusCancelled {
|
||||
continue
|
||||
}
|
||||
j.RequestCancel()
|
||||
n++
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// ActiveCount returns the number of queued or running jobs.
|
||||
func (r *Registry) ActiveCount() int {
|
||||
if r == nil {
|
||||
return 0
|
||||
}
|
||||
r.mu.RLock()
|
||||
defer r.mu.RUnlock()
|
||||
n := 0
|
||||
for _, j := range r.byID {
|
||||
if j == nil {
|
||||
continue
|
||||
}
|
||||
st := j.statusLocked()
|
||||
if st == StatusQueued || st == StatusRunning {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// Drain waits until no queued/running jobs remain or ctx is done.
|
||||
// Call RequestCancelAll first for a cooperative shutdown.
|
||||
func (r *Registry) Drain(ctx context.Context) error {
|
||||
if r == nil {
|
||||
return nil
|
||||
}
|
||||
ticker := time.NewTicker(50 * time.Millisecond)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
if r.ActiveCount() == 0 {
|
||||
return nil
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-ticker.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -21,14 +21,13 @@ func (w *Worker) runMaintenancePolicy(j *Job) {
|
||||
j.Fail("postgresql not configured")
|
||||
return
|
||||
}
|
||||
policyID, _ := j.Meta["policy_id"].(string)
|
||||
policyID = strings.TrimSpace(policyID)
|
||||
policyID := strings.TrimSpace(j.metaString("policy_id"))
|
||||
if policyID == "" {
|
||||
j.Fail("missing policy_id in job meta")
|
||||
return
|
||||
}
|
||||
dryRun, _ := j.Meta["dry_run"].(bool)
|
||||
actor, _ := j.Meta["actor_prefix"].(string)
|
||||
dryRun := j.metaBool("dry_run")
|
||||
actor := j.metaString("actor_prefix")
|
||||
ctx, cancel := j.workContext()
|
||||
defer cancel()
|
||||
|
||||
|
||||
@@ -93,9 +93,9 @@ func (w *Worker) runPostgresMaint(j *Job, kind string) {
|
||||
j.Fail("postgresql not configured")
|
||||
return
|
||||
}
|
||||
table, _ := j.Meta["table"].(string)
|
||||
dryRun, _ := j.Meta["dry_run"].(bool)
|
||||
actor, _ := j.Meta["actor_prefix"].(string)
|
||||
table := j.metaString("table")
|
||||
dryRun := j.metaBool("dry_run")
|
||||
actor := j.metaString("actor_prefix")
|
||||
ctx, cancel := j.workContext()
|
||||
defer cancel()
|
||||
auditID, _ := pgmonitor.InsertMaintenanceAudit(ctx, w.PgPool, j.TenantID, actor, kind, table, dryRun)
|
||||
|
||||
@@ -291,8 +291,8 @@ func (w *Worker) runModuleRefresh(j *Job) {
|
||||
}
|
||||
}()
|
||||
|
||||
mid, _ := j.Meta["module_id"].(string)
|
||||
if strings.TrimSpace(mid) == "" {
|
||||
mid := strings.TrimSpace(j.metaString("module_id"))
|
||||
if mid == "" {
|
||||
j.Fail("missing module_id in job meta")
|
||||
return
|
||||
}
|
||||
@@ -325,7 +325,7 @@ func (w *Worker) runTenantRefresh(j *Job) {
|
||||
}
|
||||
}()
|
||||
|
||||
moduleIDs := moduleIDsFromJobMeta(j.Meta)
|
||||
moduleIDs := moduleIDsFromJobMeta(j.metaCopy())
|
||||
if len(moduleIDs) == 0 {
|
||||
j.Fail("missing module_ids in job meta")
|
||||
return
|
||||
@@ -449,8 +449,9 @@ func (w *Worker) enqueueDeployAllSpeakers(j *Job, tenantID, revID string) {
|
||||
}
|
||||
|
||||
func (w *Worker) runDeployApply(j *Job) {
|
||||
revID, _ := j.Meta["revision_id"].(string)
|
||||
spk, hasSpeaker := j.Meta["speaker_id"].(string)
|
||||
revID := strings.TrimSpace(j.metaString("revision_id"))
|
||||
spk := strings.TrimSpace(j.metaString("speaker_id"))
|
||||
hasSpeaker := spk != ""
|
||||
if revID == "" {
|
||||
j.Fail("missing revision_id in job meta")
|
||||
return
|
||||
@@ -571,7 +572,7 @@ func (w *Worker) runDeployApply(j *Job) {
|
||||
}
|
||||
|
||||
func (w *Worker) runRollback(j *Job) {
|
||||
src, _ := j.Meta["source_revision_id"].(string)
|
||||
src := strings.TrimSpace(j.metaString("source_revision_id"))
|
||||
if src == "" {
|
||||
j.Fail("missing source_revision_id in job meta")
|
||||
return
|
||||
|
||||
Reference in New Issue
Block a user