From bc3a73063a496ebe8a3f3f3e54095ee30e7269b0 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Mon, 28 Sep 2026 10:27:09 +0100 Subject: [PATCH] Make release health-gate tests synchronize on observed state Wait for the agent supervisor to check false readiness before asserting its commit gate, and give the race-instrumented fixture a non-expiring window. Hold the monitoring request at target resolution after its history snapshot so replacement is deterministic. Keep production paths and assertions unchanged. Change-source: pulse-maintainer --- cmd/pulse-agent/main_test.go | 44 ++++++++++--- .../monitoring/metric_window_provider_test.go | 61 +++++++++++++++---- 2 files changed, 85 insertions(+), 20 deletions(-) diff --git a/cmd/pulse-agent/main_test.go b/cmd/pulse-agent/main_test.go index 146671721..2ed7a0cff 100644 --- a/cmd/pulse-agent/main_test.go +++ b/cmd/pulse-agent/main_test.go @@ -1901,11 +1901,13 @@ func (s *pendingUpdateSupervisorStub) Rollback(_ context.Context, activation age func testPendingUpdate(t *testing.T, stateDir string) *agentupdate.PendingPrivilegedUpdate { t.Helper() activation := agenthelper.UpdateResult{ - Action: "pending", - ActivationID: "pulse-agent-0123456789abcdef0123456789abcdef:0123456789abcdef", - ActiveSHA256: strings.Repeat("a", 64), - RollbackSHA256: strings.Repeat("b", 64), - RollbackDeadline: time.Now().Add(2 * time.Second).UTC(), + Action: "pending", + ActivationID: "pulse-agent-0123456789abcdef0123456789abcdef:0123456789abcdef", + ActiveSHA256: strings.Repeat("a", 64), + RollbackSHA256: strings.Repeat("b", 64), + // Leave enough time for the race-instrumented supervisor to be scheduled. + // The test below checks the health gates, not the expiry path. + RollbackDeadline: time.Now().Add(30 * time.Second).UTC(), } if err := agentupdate.PersistPendingPrivilegedUpdate(stateDir, "1.0.0", activation); err != nil { t.Fatal(err) @@ -1932,14 +1934,38 @@ func TestPendingPrivilegedUpdateCommitsOnlyAfterReadinessAndAcceptedReport(t *te reportAccepted := make(chan struct{}) close(reportAccepted) var ready atomic.Bool + notReadyChecks := make(chan struct{}, 2) + locallyReady := func() bool { + if ready.Load() { + return true + } + select { + case notReadyChecks <- struct{}{}: + default: + } + return false + } result := make(chan error, 1) go func() { - result <- supervisePendingPrivilegedUpdate(context.Background(), stub, pending, stateDir, pending.Activation.ActiveSHA256, ready.Load, reportAccepted, time.Millisecond, nil) + result <- supervisePendingPrivilegedUpdate(context.Background(), stub, pending, stateDir, pending.Activation.ActiveSHA256, locallyReady, reportAccepted, time.Millisecond, nil) }() + // Observe the supervisor checking the false readiness state on two passes. + // A sleep cannot prove it started before we flip readiness under -race. + for range 2 { + select { + case <-notReadyChecks: + case <-stub.commitCalls: + t.Fatal("pending update committed before local readiness") + case err := <-result: + t.Fatalf("pending update supervisor stopped before local readiness: %v", err) + case <-time.After(10 * time.Second): + t.Fatal("pending update supervisor did not check local readiness") + } + } select { case <-stub.commitCalls: t.Fatal("pending update committed before local readiness") - case <-time.After(20 * time.Millisecond): + default: } ready.Store(true) select { @@ -1947,7 +1973,9 @@ func TestPendingPrivilegedUpdateCommitsOnlyAfterReadinessAndAcceptedReport(t *te if activation != pending.Activation { t.Fatalf("commit activation = %#v", activation) } - case <-time.After(time.Second): + case err := <-result: + t.Fatalf("pending update supervisor stopped before committing: %v", err) + case <-time.After(10 * time.Second): t.Fatal("pending update was not committed after both health signals") } if err := <-result; err != nil { diff --git a/internal/monitoring/metric_window_provider_test.go b/internal/monitoring/metric_window_provider_test.go index 9931595f4..e29c52935 100644 --- a/internal/monitoring/metric_window_provider_test.go +++ b/internal/monitoring/metric_window_provider_test.go @@ -2,12 +2,35 @@ package monitoring import ( "strings" + "sync" "testing" "time" "github.com/rcourtman/pulse-go-rewrite/internal/alerts" + "github.com/rcourtman/pulse-go-rewrite/internal/models" + "github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources" ) +// This resolver blocks only after metricWindowPoints has captured its history. +// It lets the test replace that history at an exact point without scheduling +// sleeps or exposing a production-only hook. +type blockingMetricTargetStore struct { + entered chan struct{} + resume <-chan struct{} +} + +func (*blockingMetricTargetStore) ShouldSkipAPIPolling(string) bool { return false } +func (*blockingMetricTargetStore) GetPollingRecommendations() map[string]float64 { + return nil +} +func (*blockingMetricTargetStore) GetAll() []unifiedresources.Resource { return nil } +func (*blockingMetricTargetStore) PopulateFromSnapshot(models.StateSnapshot) {} +func (s *blockingMetricTargetStore) MetricsTargetForResource(string) *unifiedresources.MetricsTarget { + s.entered <- struct{}{} + <-s.resume + return nil +} + func TestMetricWindowPointsUsesInMemoryMetricAlias(t *testing.T) { now := time.Now().UTC() history := NewMetricsHistory(32, time.Hour) @@ -77,33 +100,47 @@ func TestMetricWindowPointsKeepsHistorySnapshotDuringReplacement(t *testing.T) { }, } - history.metricWindowMu.Lock() - result := make(chan []alerts.MetricWindowPoint, 1) + resume := make(chan struct{}) + var resumeOnce sync.Once + resumeRequest := func() { resumeOnce.Do(func() { close(resume) }) } + defer resumeRequest() + entered := make(chan struct{}, 1) + monitor.resourceStore = &blockingMetricTargetStore{entered: entered, resume: resume} + type requestResult struct { + points []alerts.MetricWindowPoint + err error + } + result := make(chan requestResult, 1) go func() { - points, _ := monitor.metricWindowPoints(alerts.MetricWindowRequest{ + points, err := monitor.metricWindowPoints(alerts.MetricWindowRequest{ ResourceID: "vm-1", ResourceType: "vm", Metric: "cpu", Start: now.Add(-5 * time.Minute), End: now, }) - result <- points + result <- requestResult{points: points, err: err} }() - // Give the request time to snapshot history and block on its cache mutex, - // matching the mock seed replacement that exposed the startup panic. - time.Sleep(20 * time.Millisecond) + select { + case <-entered: + case <-time.After(10 * time.Second): + t.Fatal("metric request did not reach target resolution after snapshot") + } monitor.mu.Lock() monitor.metricsHistory = NewMetricsHistory(32, time.Hour) monitor.mu.Unlock() - history.metricWindowMu.Unlock() + resumeRequest() select { - case points := <-result: - if len(points) != 1 || points[0].Value != 42 { - t.Fatalf("metricWindowPoints = %+v, want cached point from the request snapshot", points) + case got := <-result: + if got.err != nil { + t.Fatalf("metricWindowPoints returned error: %v", got.err) } - case <-time.After(time.Second): + if len(got.points) != 1 || got.points[0].Value != 42 { + t.Fatalf("metricWindowPoints = %+v, want cached point from the request snapshot", got.points) + } + case <-time.After(10 * time.Second): t.Fatal("metricWindowPoints did not finish after history replacement") } }