fix(darwin): probe interception when stabilization finishes

The post-stabilization reconcile verifies rule text, which cannot tell a
live redirect from an anchor pf has stopped evaluating. Sleep/wake QA
caught exactly that split: references intact, anchor rules intact,
post-load verification passed, and every query through the system
resolver timing out while the direct listener answered.

Nothing else probed. The interception probe monitor stands down while
stabilization owns pf and is never re-armed afterwards, so functional
recovery waited for the periodic watchdog - 11 seconds in the captured
run, up to a full 30-second interval - on a host whose link and default
route were already back. The watchdog's probe then failed once, forced a
reload, and public and VPN split-DNS both recovered immediately.

Probe once at the end of stabilization and, if it fails, force exactly
one reload and confirm with one more probe. Not the probe monitor: that
keeps probing for ~7.5s and can force a reload per failed probe, where
this path needs a single bounded repair before handing back to the
watchdog. Skipped when a monitor already owns probing, when intercept
state is gone, or during exec backoff.

Hand ownership over deterministically rather than skipping on sight. A
probe monitor started by an ignored network change claimed
functional-probe ownership before checking whether it could work, then
stood down because stabilization still owned pf; the verifier read that
claimed flag as "somebody is probing" and skipped, so neither path
probed and recovery fell back to the watchdog anyway. The monitor now
checks eligibility before claiming, and the verifier waits out a holder
that releases, yielding only to one that keeps probing.

Extract the completion block into finishPFStabilization so the wiring is
testable, and cover the bounded repair, the healthy path that must not
reload, a prober that claims and stands down, a prober that keeps
working, and a monitor that must not claim ownership while stabilizing.
This commit is contained in:
Cuong Manh Le
2026-08-21 14:47:58 +07:00
parent a828c8853a
commit d7c30b18ed
2 changed files with 303 additions and 18 deletions
+123 -18
View File
@@ -1403,25 +1403,123 @@ func (p *prog) pfStabilizationLoopWithMaxWait(ctx context.Context, stableRequire
}
if time.Since(stableSince) >= stableRequired {
// The active loop retains ownership until the dedicated post-stable
// repair finishes. It never re-enters stabilization recursively.
mainLog.Load().Info().Msgf("DNS intercept: pf stable for %s — reconciling anchor rules", stableRequired)
result := p.reconcilePFAnchorAfterStabilization()
if result != pfAnchorCheckRestored && result != pfAnchorCheckIntact {
p.scheduleDelayedRechecks()
}
routes, domainlessServers, exemptions := p.refreshDNSAfterVPNSettle("pf_stabilized")
if routes == 0 && domainlessServers == 0 && exemptions == 0 {
p.scheduleDNSAfterVPNSettleRefresh("pf_stabilized_followup", pfAnchorRecheckDelayLong)
}
if p.hasPendingTunnelReconcile() {
p.scheduleDelayedRechecks()
}
p.finishPFStabilization(stableRequired)
return
}
}
}
// finishPFStabilization runs the work stabilization exists to do, once the ruleset has
// held still for the required window. The active loop retains ownership throughout: it
// never re-enters stabilization recursively.
func (p *prog) finishPFStabilization(stableRequired time.Duration) {
mainLog.Load().Info().Msgf("DNS intercept: pf stable for %s — reconciling anchor rules", stableRequired)
result := p.reconcilePFAnchorAfterStabilization()
if result != pfAnchorCheckRestored && result != pfAnchorCheckIntact {
p.scheduleDelayedRechecks()
}
routes, domainlessServers, exemptions := p.refreshDNSAfterVPNSettle("pf_stabilized")
if routes == 0 && domainlessServers == 0 && exemptions == 0 {
p.scheduleDNSAfterVPNSettleRefresh("pf_stabilized_followup", pfAnchorRecheckDelayLong)
}
if p.hasPendingTunnelReconcile() {
p.scheduleDelayedRechecks()
}
p.verifyInterceptAfterStabilization()
}
// probePFInterceptFn and forceReloadPFInterceptFn are the functional verification seams
// shared by both probers.
var (
probePFInterceptFn = (*prog).probePFIntercept
forceReloadPFInterceptFn = (*prog).forceReloadPFMainRuleset
)
// pfFunctionalProbeOwnerWait bounds how long the post-stabilization verifier waits for
// another prober to release functional-probe ownership. A probe monitor that is only
// standing down releases it at once; one that is genuinely probing holds it for its whole
// window, and then the verifier steps aside. A var so tests can shorten the wait.
var pfFunctionalProbeOwnerWait = 2 * time.Second
// pfFunctionalProbeOwnerPoll is how often that wait re-tries the claim.
const pfFunctionalProbeOwnerPoll = 25 * time.Millisecond
// interceptProbeMonitorAllowed reports whether the probe monitor may run at all.
//
// The monitor must consult this before claiming ownership. Claiming first and checking
// second means a monitor that is about to stand down still takes the flag, and the
// post-stabilization verifier - which sees a set flag as "somebody else is probing" -
// skips. Neither probes, and the outage lasts until the next watchdog tick.
func (p *prog) interceptProbeMonitorAllowed() bool {
return p.dnsInterceptState != nil && !p.pfStabilizing.Load()
}
// claimFunctionalProbeOwner takes ownership of functional probing, waiting up to wait for
// a current owner to release it. It reports whether ownership was acquired; the caller
// releases with pfMonitorRunning.Store(false).
func (p *prog) claimFunctionalProbeOwner(wait time.Duration) bool {
deadline := time.Now().Add(wait)
for {
if p.pfMonitorRunning.CompareAndSwap(false, true) {
return true
}
if !time.Now().Before(deadline) {
return false
}
time.Sleep(pfFunctionalProbeOwnerPoll)
}
}
// verifyInterceptAfterStabilization proves pf is actually translating once stabilization
// has finished, and repairs it once if it is not.
//
// The reconcile above verifies rule text. That cannot distinguish a live redirect from an
// anchor pf has stopped evaluating, and after sleep/wake with a VPN reconnect those come
// apart: rules present, references present, post-load verification passed, and every query
// through the system resolver timing out. Nothing else notices until the periodic watchdog
// runs its own probe - the interception probe monitor stands down while stabilization owns
// pf and is not re-armed afterwards - so recovery waits up to a full watchdog interval on a
// host whose link and default route are already back.
//
// One probe, then at most one forced reload and one confirming probe. Deliberately not the
// probe monitor: that keeps probing for ~7.5s and can force a reload per failed probe,
// where this path needs a single bounded repair and then hands back to the watchdog.
func (p *prog) verifyInterceptAfterStabilization() {
if p.dnsInterceptState == nil || p.pfExecBackoffActive() {
return
}
// The probe monitor is the other functional prober, so only one of us may run - but
// "somebody holds the flag" is not the same as "somebody is probing". A monitor that
// started while stabilization owns pf stands down immediately, and skipping on sight
// left nobody probing at all. Wait briefly for the holder to release instead.
if !p.claimFunctionalProbeOwner(pfFunctionalProbeOwnerWait) {
mainLog.Load().Warn().Msgf("DNS intercept: post-stabilization probe skipped — another prober held ownership for %s; leaving recovery to the watchdog", pfFunctionalProbeOwnerWait)
return
}
defer p.pfMonitorRunning.Store(false)
// Ownership can take a moment to arrive; make sure there is still an intercept to check.
if p.dnsInterceptState == nil {
return
}
if probePFInterceptFn(p) {
mainLog.Load().Debug().Msg("DNS intercept: post-stabilization probe passed — interception is translating")
return
}
mainLog.Load().Warn().Msg("DNS intercept: post-stabilization rules are intact but the probe FAILED — forcing one reload")
if !forceReloadPFInterceptFn(p) {
mainLog.Load().Error().Msg("DNS intercept: post-stabilization forced reload did not run — leaving recovery to the watchdog")
return
}
if probePFInterceptFn(p) {
mainLog.Load().Info().Msg("DNS intercept: interception restored by the post-stabilization reload")
return
}
mainLog.Load().Error().Msg("DNS intercept: interception still not translating after the post-stabilization reload — the watchdog will retry")
}
var runPFAnchorCheckCommand = func(args ...string) ([]byte, error) {
return exec.Command("pfctl", args...).CombinedOutput()
}
@@ -1930,6 +2028,13 @@ func buildDNSQueryPacket(domain string) []byte {
// The backoff schedule provides both fast detection (immediate + 500ms) and extended
// coverage (up to ~8s) to win the race against async pf reloads by hypervisors.
func (p *prog) pfInterceptMonitor() {
// Eligibility first, ownership second. A monitor that is about to stand down must not
// take the flag on its way out: the post-stabilization verifier reads that flag as
// "another prober is working" and would step aside for a prober that never probes.
if !p.interceptProbeMonitorAllowed() {
mainLog.Load().Debug().Msg("DNS intercept monitor: not starting — intercept disabled or stabilizing")
return
}
if !p.pfMonitorRunning.CompareAndSwap(false, true) {
mainLog.Load().Debug().Msg("DNS intercept monitor: already running, skipping")
return
@@ -1946,23 +2051,23 @@ func (p *prog) pfInterceptMonitor() {
if delay > 0 {
time.Sleep(delay)
}
if p.dnsInterceptState == nil || p.pfStabilizing.Load() {
if !p.interceptProbeMonitorAllowed() {
mainLog.Load().Debug().Msg("DNS intercept monitor: aborting — intercept disabled or stabilizing")
return
}
if p.probePFIntercept() {
if probePFInterceptFn(p) {
mainLog.Load().Debug().Msgf("DNS intercept monitor: probe %d/%d passed", i+1, len(delays))
continue // working now — keep monitoring in case it breaks later in the window
}
// Probe failed — pf translation is broken. Force full reload.
mainLog.Load().Warn().Msgf("DNS intercept monitor: probe %d/%d FAILED — pf translation broken, forcing full ruleset reload", i+1, len(delays))
p.forceReloadPFMainRuleset()
forceReloadPFInterceptFn(p)
// Verify the reload fixed it
time.Sleep(200 * time.Millisecond)
if p.probePFIntercept() {
if probePFInterceptFn(p) {
mainLog.Load().Info().Msg("DNS intercept monitor: probe passed after reload — interception restored")
// Continue monitoring in case the hypervisor reloads pf again
} else {
+180
View File
@@ -828,3 +828,183 @@ func TestExemptVPNDNSServersDeferredWhileStabilizing(t *testing.T) {
t.Error("pfEnsureRunning was left held by a deferred exemption")
}
}
// stubStabilizationProbe replaces the post-stabilization verification seams and returns
// counters for probe and forced-reload calls.
func stubStabilizationProbe(t *testing.T, probeResults []bool, reloadOK bool) (probes, reloads *int) {
t.Helper()
originalProbe, originalReload := probePFInterceptFn, forceReloadPFInterceptFn
t.Cleanup(func() {
probePFInterceptFn, forceReloadPFInterceptFn = originalProbe, originalReload
})
probeCalls, reloadCalls := 0, 0
probePFInterceptFn = func(*prog) bool {
result := false
if probeCalls < len(probeResults) {
result = probeResults[probeCalls]
}
probeCalls++
return result
}
forceReloadPFInterceptFn = func(*prog) bool {
reloadCalls++
return reloadOK
}
return &probeCalls, &reloadCalls
}
// TestPostStabilizationVerifiesInterceptionFunctionally is the post-wake continuity
// boundary: the reconcile above it only proves rule text, and QA saw rules intact,
// references intact and post-load verification passed while every query through the system
// resolver timed out. Nothing else probes until the periodic watchdog, because the probe
// monitor stands down while stabilization owns pf, so recovery waited for that tick.
func TestPostStabilizationVerifiesInterceptionFunctionally(t *testing.T) {
// Probe fails once, then passes after the reload.
probes, reloads := stubStabilizationProbe(t, []bool{false, true}, true)
p := &prog{dnsInterceptState: &pfState{}}
p.verifyInterceptAfterStabilization()
if *probes != 2 {
t.Errorf("probe calls = %d, want 2: one to detect and one to confirm the repair", *probes)
}
if *reloads != 1 {
t.Errorf("forced reloads = %d, want exactly 1 bounded repair", *reloads)
}
}
// TestPostStabilizationProbePassSkipsReload keeps the healthy path free of a pf reload,
// which would flush states and kill in-flight DoH connections for nothing.
func TestPostStabilizationProbePassSkipsReload(t *testing.T) {
probes, reloads := stubStabilizationProbe(t, []bool{true}, true)
p := &prog{dnsInterceptState: &pfState{}}
p.verifyInterceptAfterStabilization()
if *probes != 1 || *reloads != 0 {
t.Errorf("probe calls = %d, forced reloads = %d, want 1/0", *probes, *reloads)
}
}
// TestPostStabilizationRepairIsBounded pins the "one bounded recovery" contract: a probe
// that never passes must not turn into a reload loop here - the watchdog owns retries.
func TestPostStabilizationRepairIsBounded(t *testing.T) {
probes, reloads := stubStabilizationProbe(t, []bool{false, false, false}, true)
p := &prog{dnsInterceptState: &pfState{}}
p.verifyInterceptAfterStabilization()
if *reloads != 1 {
t.Errorf("forced reloads = %d, want 1: the repair must not loop", *reloads)
}
if *probes != 2 {
t.Errorf("probe calls = %d, want 2", *probes)
}
}
// TestPostStabilizationWaitsForAProberThatStandsDown is the interleaving that made
// "skip when the flag is set" wrong. A probe monitor started by an ignored network change
// claims functional-probe ownership and then aborts, because stabilization still owns pf.
// If the verifier treats the claimed flag as "somebody is probing", neither path probes and
// the outage lasts until the next watchdog tick - the exact window this is meant to close.
func TestPostStabilizationWaitsForAProberThatStandsDown(t *testing.T) {
probes, reloads := stubStabilizationProbe(t, []bool{false, true}, true)
p := &prog{dnsInterceptState: &pfState{}}
// Model the monitor's claim-then-abort: ownership is held, then released.
p.pfMonitorRunning.Store(true)
released := make(chan struct{})
go func() {
time.Sleep(50 * time.Millisecond)
p.pfMonitorRunning.Store(false)
close(released)
}()
p.verifyInterceptAfterStabilization()
<-released
if *probes != 2 {
t.Errorf("probe calls = %d, want 2: the verifier must wait out a prober that stands down", *probes)
}
if *reloads != 1 {
t.Errorf("forced reloads = %d, want 1", *reloads)
}
}
// TestPostStabilizationYieldsToAProberThatKeepsProbing is the other half of the handoff:
// when the holder is genuinely working through its probe sequence, the verifier must step
// aside rather than run a second prober against the same pf state.
func TestPostStabilizationYieldsToAProberThatKeepsProbing(t *testing.T) {
originalWait := pfFunctionalProbeOwnerWait
pfFunctionalProbeOwnerWait = 30 * time.Millisecond
t.Cleanup(func() { pfFunctionalProbeOwnerWait = originalWait })
probes, reloads := stubStabilizationProbe(t, []bool{false}, true)
p := &prog{dnsInterceptState: &pfState{}}
p.pfMonitorRunning.Store(true) // held for the whole wait
p.verifyInterceptAfterStabilization()
if *probes != 0 || *reloads != 0 {
t.Errorf("probe calls = %d, forced reloads = %d, want 0/0 while another prober is working", *probes, *reloads)
}
}
// TestInterceptMonitorDoesNotClaimOwnershipWhileStabilizing pins the source of that race:
// a monitor which cannot do useful work must not take functional-probe ownership on its way
// out, or it starves the post-stabilization verifier.
func TestInterceptMonitorDoesNotClaimOwnershipWhileStabilizing(t *testing.T) {
probes, reloads := stubStabilizationProbe(t, []bool{false}, true)
p := &prog{dnsInterceptState: &pfState{}}
p.pfStabilizing.Store(true)
if p.interceptProbeMonitorAllowed() {
t.Fatal("the probe monitor considers itself eligible while stabilization owns pf")
}
p.pfInterceptMonitor()
if *probes != 0 || *reloads != 0 {
t.Errorf("probe calls = %d, forced reloads = %d, want 0/0 from a monitor that cannot run", *probes, *reloads)
}
if !p.claimFunctionalProbeOwner(0) {
t.Error("the aborted monitor left functional-probe ownership taken; the verifier would skip")
}
p.pfMonitorRunning.Store(false)
}
// TestFinishPFStabilizationRunsFunctionalVerification wires the fix to the production
// completion path: deleting the verification call, or reordering it before the reconcile,
// makes this fail.
func TestFinishPFStabilizationRunsFunctionalVerification(t *testing.T) {
stubPFAnchorCheckCommand(t, map[string]string{
"-sn": `rdr-anchor "com.controld.ctrld"`,
"-sr": `anchor "com.controld.ctrld"`,
"-a com.controld.ctrld -sr": "pass in quick on lo0",
"-a com.controld.ctrld -sn": "rdr on lo0",
})
originalResolver := initializeOsResolver
initializeOsResolver = func(bool) []string { return nil }
t.Cleanup(func() { initializeOsResolver = originalResolver })
probes, reloads := stubStabilizationProbe(t, []bool{false, true}, true)
p := &prog{dnsInterceptState: &pfState{}}
p.pfStabilizing.Store(true)
p.finishPFStabilization(time.Millisecond)
if *probes == 0 {
t.Fatal("stabilization completed without probing functional interception; recovery would wait for the watchdog")
}
if *reloads != 1 {
t.Errorf("forced reloads = %d, want 1", *reloads)
}
p.pfDelayedRecheckMu.Lock()
timers := append([]*time.Timer(nil), p.pfDelayedRecheckTimers...)
p.pfDelayedRecheckTimers = nil
p.pfDelayedRecheckMu.Unlock()
for _, timer := range timers {
timer.Stop()
}
}