mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-02 20:29:43 +00:00
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
This commit is contained in:
parent
5b3a47e305
commit
bc3a73063a
2 changed files with 85 additions and 20 deletions
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue