db: reap ephemeral GC watcher goroutine on cancel and reschedule

A stopped timer never fires, so the per-node watcher leaked until close.
This commit is contained in:
Kristoffer Dalby
2026-06-05 14:05:56 +00:00
committed by Kristoffer Dalby
parent 06d6816dc9
commit 4c165ae5e7
2 changed files with 58 additions and 10 deletions
+28 -10
View File
@@ -393,17 +393,25 @@ type EphemeralGarbageCollector struct {
mu sync.Mutex
deleteFunc func(types.NodeID)
toBeDeleted map[types.NodeID]*time.Timer
toBeDeleted map[types.NodeID]ephemeralTimer
deleteCh chan types.NodeID
cancelCh chan struct{}
}
// ephemeralTimer pairs a node's pending-deletion timer with a done channel
// used to reap its watcher goroutine on Cancel or reschedule. Without it a
// stopped timer never fires and the goroutine leaks until Close.
type ephemeralTimer struct {
timer *time.Timer
done chan struct{}
}
// NewEphemeralGarbageCollector creates a new [EphemeralGarbageCollector], it takes
// a deleteFunc that will be called when a node is scheduled for deletion.
func NewEphemeralGarbageCollector(deleteFunc func(types.NodeID)) *EphemeralGarbageCollector {
return &EphemeralGarbageCollector{
toBeDeleted: make(map[types.NodeID]*time.Timer),
toBeDeleted: make(map[types.NodeID]ephemeralTimer),
deleteCh: make(chan types.NodeID, 10),
cancelCh: make(chan struct{}),
deleteFunc: deleteFunc,
@@ -416,8 +424,8 @@ func (e *EphemeralGarbageCollector) Close() {
defer e.mu.Unlock()
// Stop all timers
for _, timer := range e.toBeDeleted {
timer.Stop()
for _, t := range e.toBeDeleted {
t.timer.Stop()
}
// Close the cancel channel to signal all goroutines to exit
@@ -440,13 +448,16 @@ func (e *EphemeralGarbageCollector) Schedule(nodeID types.NodeID, expiry time.Du
// Continue with scheduling
}
// If a timer already exists for this node, stop it first
if oldTimer, exists := e.toBeDeleted[nodeID]; exists {
oldTimer.Stop()
// If a timer already exists for this node, stop it and reap its
// watcher goroutine before scheduling a fresh one.
if old, exists := e.toBeDeleted[nodeID]; exists {
old.timer.Stop()
close(old.done)
}
timer := time.NewTimer(expiry)
e.toBeDeleted[nodeID] = timer
done := make(chan struct{})
e.toBeDeleted[nodeID] = ephemeralTimer{timer: timer, done: done}
// Start a goroutine to handle the timer completion
go func() {
select {
@@ -461,7 +472,13 @@ func (e *EphemeralGarbageCollector) Schedule(nodeID types.NodeID, expiry time.Du
case <-e.cancelCh:
// GC is shutting down, don't send to deleteCh
return
case <-done:
// Cancelled or rescheduled before the send landed.
return
}
case <-done:
// Cancelled or rescheduled before the timer fired.
return
case <-e.cancelCh:
// If the GC is closed, exit the goroutine
return
@@ -474,8 +491,9 @@ func (e *EphemeralGarbageCollector) Cancel(nodeID types.NodeID) {
e.mu.Lock()
defer e.mu.Unlock()
if timer, ok := e.toBeDeleted[nodeID]; ok {
timer.Stop()
if t, ok := e.toBeDeleted[nodeID]; ok {
t.timer.Stop()
close(t.done)
delete(e.toBeDeleted, nodeID)
}
}