diff --git a/internal/jobs/worker.go b/internal/jobs/worker.go index f1e5bb8..e1cecc5 100644 --- a/internal/jobs/worker.go +++ b/internal/jobs/worker.go @@ -100,22 +100,11 @@ func (w *Worker) Process(j *Job) { j.Fail("missing module_id in job meta") return } - rev, err := pipeline.RefreshModule(context.Background(), w.Store, w.httpClient(), j.TenantID, mid) - if err != nil { + if err := pipeline.RefreshModuleIngest(context.Background(), w.Store, w.httpClient(), j.TenantID, mid); err != nil { j.Fail(err.Error()) return } - j.mergeMeta(map[string]any{"revision_id": rev}) - if entries, total, err := w.buildRevisionLogEntries(j.TenantID, rev); err == nil { - j.mergeMeta(map[string]any{ - "log_entries": entries, - "log_total": total, - "log_generated": time.Now().UTC().Format(time.RFC3339Nano), - }) - } else { - j.mergeMeta(map[string]any{"log_build_error": err.Error()}) - } - w.finishModuleRefreshSuccess(j, rev) + w.finishModuleRefreshSuccess(j, mid) case KindDeployApply: w.runDeployApply(j) case KindRevisionRollback: @@ -146,27 +135,46 @@ func (w *Worker) tenantRefreshMu(tenantID string) *sync.Mutex { return v.(*sync.Mutex) } -// finishModuleRefreshSuccess marks the job succeeded and enqueues deploy_apply only when no other -// module_refresh is still queued or running for the same tenant (coalesces parallel refreshes). -func (w *Worker) finishModuleRefreshSuccess(j *Job, rev string) { - if w == nil || w.Registry == nil { +// finishModuleRefreshSuccess marks the refresh job and, for the last active refresh in tenant, +// creates one aggregate revision and enqueues a single deploy_apply. +func (w *Worker) finishModuleRefreshSuccess(j *Job, triggerModuleID string) { + if w == nil || w.Store == nil { j.Succeed() return } mu := w.tenantRefreshMu(j.TenantID) mu.Lock() - deferDeploy := w.Registry.CountOtherActiveModuleRefresh(j.TenantID, j.ID) > 0 + defer mu.Unlock() + deferDeploy := false + if w.Registry != nil { + deferDeploy = w.Registry.CountOtherActiveModuleRefresh(j.TenantID, j.ID) > 0 + } if deferDeploy { j.mergeMeta(map[string]any{ "deploy_apply_deferred": true, "deploy_apply_defer_reason": "parallel_module_refresh", }) + j.Succeed() + return + } + + rev, err := pipeline.RenderTenantRevision(context.Background(), w.Store, w.httpClient(), j.TenantID, triggerModuleID) + if err != nil { + j.Fail(err.Error()) + return + } + j.mergeMeta(map[string]any{"revision_id": rev}) + if entries, total, err := w.buildRevisionLogEntries(j.TenantID, rev); err == nil { + j.mergeMeta(map[string]any{ + "log_entries": entries, + "log_total": total, + "log_generated": time.Now().UTC().Format(time.RFC3339Nano), + }) + } else { + j.mergeMeta(map[string]any{"log_build_error": err.Error()}) } j.Succeed() - mu.Unlock() - if !deferDeploy { - w.enqueueDeployAllSpeakers(j, j.TenantID, rev) - } + w.enqueueDeployAllSpeakers(j, j.TenantID, rev) } // enqueueDeployAllSpeakers queues the same work as POST /v1/apply (all speakers, no speaker_id). diff --git a/internal/pipeline/refresh.go b/internal/pipeline/refresh.go index 78854e2..c54721b 100644 --- a/internal/pipeline/refresh.go +++ b/internal/pipeline/refresh.go @@ -37,28 +37,34 @@ func MaterializedASPrefixKey(asn int64) string { return fmt.Sprintf("as:%d", asn) } -// RefreshModule runs ingest (where applicable) for one module, then renders a new revision whose -// BIRD materialization includes prefixes from all enabled modules of the tenant (others via live collect). -// If the tenant-wide materialized prefix set is unchanged from the latest revision, returns that -// revision id and does not insert a duplicate config_revision. -func RefreshModule(ctx context.Context, st store.Backend, hc *http.Client, tenantID, moduleID string) (revisionID string, err error) { +// RefreshModuleIngest runs ingest for one module and persists side-effects (ASN metadata, CDN etags, etc). +// It does not create a new config revision. +func RefreshModuleIngest(ctx context.Context, st store.Backend, hc *http.Client, tenantID, moduleID string) error { if hc == nil { hc = http.DefaultClient } mod, err := st.GetModule(tenantID, moduleID) if err != nil { - return "", err + return err } if !mod.Enabled { - return "", fmt.Errorf("module disabled") + return fmt.Errorf("module disabled") } - rows, err := collectModulePrefixRows(ctx, st, hc, tenantID, mod) + _, err = collectModulePrefixRows(ctx, st, hc, tenantID, mod) if err != nil { - return "", err + return err } + return nil +} - agg, err := aggregateTenantPrefixRows(ctx, st, hc, tenantID, moduleID, rows) +// RenderTenantRevision renders one tenant-wide revision using current data from all enabled modules. +// If materialized prefixes are unchanged, returns latest revision id without creating a duplicate. +func RenderTenantRevision(ctx context.Context, st store.Backend, hc *http.Client, tenantID, triggerModuleID string) (revisionID string, err error) { + if hc == nil { + hc = http.DefaultClient + } + agg, err := aggregateTenantPrefixRowsAll(ctx, st, hc, tenantID) if err != nil { return "", err } @@ -69,17 +75,25 @@ func RefreshModule(ctx context.Context, st store.Backend, hc *http.Client, tenan agg = smartAggregatePrefixRows(agg) revisionID = uuid.NewString() - parent := parentRevision(st, tenantID, moduleID) - preview, err := buildPreviewFragments(st, tenantID, moduleID, revisionID, agg) + parent := parentRevision(st, tenantID, triggerModuleID) + preview, err := buildPreviewFragments(st, tenantID, triggerModuleID, revisionID, agg) if err != nil { return "", err } - if err := st.CreateRenderRevision(revisionID, tenantID, moduleID, parent, hash, preview, agg); err != nil { + if err := st.CreateRenderRevision(revisionID, tenantID, triggerModuleID, parent, hash, preview, agg); err != nil { return "", err } return revisionID, nil } +// RefreshModule keeps backwards-compatible behavior: module ingest + immediate tenant render. +func RefreshModule(ctx context.Context, st store.Backend, hc *http.Client, tenantID, moduleID string) (revisionID string, err error) { + if err := RefreshModuleIngest(ctx, st, hc, tenantID, moduleID); err != nil { + return "", err + } + return RenderTenantRevision(ctx, st, hc, tenantID, moduleID) +} + // collectModulePrefixRows returns materialized prefix rows for a single module (source of truth from store / ASN resolve / CDN fetch). func collectModulePrefixRows(ctx context.Context, st store.Backend, hc *http.Client, tenantID string, mod *store.Module) ([]store.PrefixRow, error) { moduleID := mod.ID @@ -452,21 +466,15 @@ func ipToHostPrefix(ip netip.Addr) string { return netip.PrefixFrom(ip, bits).Masked().String() } -// aggregateTenantPrefixRows builds the union of materialized prefixes for all enabled modules. -// The module that triggered refresh contributes freshRows; every other module is collected live from the store -// (same logic as refresh). We do not reuse other modules' saved revisions as prefix sources, because each revision -// already stores the full tenant-wide aggregate — mixing them with freshRows would duplicate prefixes. -func aggregateTenantPrefixRows(ctx context.Context, st store.Backend, hc *http.Client, tenantID, changedModuleID string, freshRows []store.PrefixRow) ([]store.PrefixRow, error) { +// aggregateTenantPrefixRowsAll builds the union of materialized prefixes for all enabled modules +// using current source data from store/external resolvers. +func aggregateTenantPrefixRowsAll(ctx context.Context, st store.Backend, hc *http.Client, tenantID string) ([]store.PrefixRow, error) { mods := st.ListModules(tenantID) var out []store.PrefixRow for _, m := range mods { if m == nil || !m.Enabled { continue } - if m.ID == changedModuleID { - out = append(out, freshRows...) - continue - } omod, err := st.GetModule(tenantID, m.ID) if err != nil { return nil, err