Merge current Pulse upstream for publication

Incorporate the metrics startup hook capture from PR #1868 while retaining
the reviewed causal cleanup and barrier-ordering coverage. This advances the
open publication proposal without rewriting any accepted commit.

Change-source: pulse-maintainer
Contract-Neutral: Integration reconciliation only; no additional public contract delta.
This commit is contained in:
pulse-triage[bot] 2026-09-02 16:09:57 +01:00
commit 111b1bd251
2 changed files with 21 additions and 21 deletions

View file

@ -211,6 +211,10 @@ type Store struct {
identityMigrationPending atomic.Bool
commercialRetentionSeconds atomic.Int64
commercialPurgeEligibleAt atomic.Int64
// startupHook is captured once at construction so a store that outlives
// the test that built it never invokes a later test's hook.
startupHook func()
}
// SetCommercialHistoryRetention applies a delayed commercial ceiling to the
@ -303,6 +307,7 @@ func NewStore(config StoreConfig) (*Store, error) {
stopCh: make(chan struct{}),
doneCh: make(chan struct{}),
maintenanceDoneCh: make(chan struct{}),
startupHook: startupMaintenanceHook,
}
// Initialize schema
@ -883,8 +888,8 @@ func (s *Store) WaitForMaintenance(timeout time.Duration) error {
func (s *Store) runStartupMaintenance() {
start := time.Now()
if startupMaintenanceHook != nil {
startupMaintenanceHook()
if s.startupHook != nil {
s.startupHook()
}
if s.identityMigrationPending.Swap(false) {

View file

@ -618,12 +618,11 @@ func TestStoreFlushMakesQueuedWritesVisible(t *testing.T) {
func TestNewStoreDefersStartupMaintenance(t *testing.T) {
previousHook := startupMaintenanceHook
// The hook parks the maintenance worker before any startup work runs, so
// startup maintenance cannot complete until the test releases it.
started := make(chan struct{})
release := make(chan struct{})
var releaseOnce sync.Once
releaseMaintenance := func() {
releaseOnce.Do(func() { close(release) })
}
startupMaintenanceHook = func() {
close(started)
<-release
@ -642,11 +641,11 @@ func TestNewStoreDefersStartupMaintenance(t *testing.T) {
store, err = NewStore(cfg)
close(done)
}()
// Release the worker and wait for NewStore before restoring the hook, so a
// failed run never leaks a parked store into the next test.
t.Cleanup(func() {
// Always unblock and join the worker before restoring the process-wide
// hook. Otherwise a failed assertion can leak this worker into the next
// test, where it may invoke that test's hook.
releaseMaintenance()
releaseOnce.Do(func() { close(release) })
<-done
if store != nil {
store.Close()
@ -654,21 +653,17 @@ func TestNewStoreDefersStartupMaintenance(t *testing.T) {
startupMaintenanceHook = previousHook
})
select {
case <-started:
case <-time.After(5 * time.Second):
t.Fatal("startup maintenance was not scheduled")
}
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("NewStore blocked on startup maintenance")
}
// Deferral is proven by ordering, not by a wall-clock bound: NewStore must
// return while the worker is still parked in the hook. If startup
// maintenance ever ran inline again, this receive would block until the
// package test timeout instead of flaking on a slow CI disk.
<-started
<-done
if err != nil {
t.Fatalf("NewStore returned error: %v", err)
}
releaseOnce.Do(func() { close(release) })
}
func TestStoreWaitForMaintenanceWaitsForQueuedStartupWork(t *testing.T) {