From c07f28cd047b8550eb4b7224f7903f88c7da3372 Mon Sep 17 00:00:00 2001 From: Nate Brown Date: Thu, 23 Jul 2026 15:24:22 -0500 Subject: [PATCH] Fold the rebind counter and traffic flags into one atomic word --- connection_manager.go | 21 ++++++----- connection_manager_test.go | 72 ++++++++++++++++++------------------- control.go | 2 +- control_tester.go | 10 ++++++ e2e/rebind_test.go | 54 ++++++++++++++++++++++++++++ hostmap.go | 73 +++++++++++++++++++++++++++++++++----- hostmap_test.go | 40 +++++++++++++++++++++ inside.go | 14 +++----- interface.go | 4 +-- 9 files changed, 225 insertions(+), 65 deletions(-) diff --git a/connection_manager.go b/connection_manager.go index 88f31321..26bd4d42 100644 --- a/connection_manager.go +++ b/connection_manager.go @@ -105,11 +105,17 @@ func (cm *connectionManager) getInactivityTimeout() time.Duration { } func (cm *connectionManager) In(h *HostInfo) { - h.in.Store(true) + h.markIn() } -func (cm *connectionManager) Out(h *HostInfo) { - h.out.Store(true) +// OutRelay records relayed traffic, leaving the rebind epoch for the direct path to this host to consume +func (cm *connectionManager) OutRelay(h *HostInfo) { + h.markOutOnly() +} + +// Out records outbound traffic and reports whether we rebound since this tunnel last sent +func (cm *connectionManager) Out(h *HostInfo) bool { + return h.markOut(cm.intf.rebindEpoch.Load()) } func (cm *connectionManager) RelayUsed(localIndex uint32) { @@ -128,8 +134,7 @@ func (cm *connectionManager) RelayUsed(localIndex uint32) { // getAndResetTrafficCheck returns if there was any inbound or outbound traffic within the last tick and // resets the state for this local index func (cm *connectionManager) getAndResetTrafficCheck(h *HostInfo, now time.Time) (bool, bool) { - in := h.in.Swap(false) - out := h.out.Swap(false) + in, out := h.takeTraffic() if in || out { h.lastUsed = now } @@ -340,7 +345,7 @@ func (cm *connectionManager) makeTrafficDecision(localIndex uint32, now time.Tim "tunnelCheck", m{"state": "alive", "method": "passive"}, ) } - hostinfo.pendingDeletion.Store(false) + hostinfo.setPendingDeletion(false) if mainHostInfo { decision = tryRehandshake @@ -363,7 +368,7 @@ func (cm *connectionManager) makeTrafficDecision(localIndex uint32, now time.Tim return decision, hostinfo, primary } - if hostinfo.pendingDeletion.Load() { + if hostinfo.isPendingDeletion() { // We have already sent a test packet and nothing was returned, this hostinfo is dead hostinfo.logger(cm.l).Info("Tunnel status", "tunnelCheck", m{"state": "dead", "method": "active"}, @@ -414,7 +419,7 @@ func (cm *connectionManager) makeTrafficDecision(localIndex uint32, now time.Tim } } - hostinfo.pendingDeletion.Store(true) + hostinfo.setPendingDeletion(true) cm.trafficTimer.Add(hostinfo.localIndexId, cm.pendingDeletionInterval) return decision, hostinfo, nil } diff --git a/connection_manager_test.go b/connection_manager_test.go index 25637c25..09d4c88c 100644 --- a/connection_manager_test.go +++ b/connection_manager_test.go @@ -86,25 +86,25 @@ func Test_NewConnectionManagerTest(t *testing.T) { // We saw traffic out to vpnIp nc.Out(hostinfo) nc.In(hostinfo) - assert.False(t, hostinfo.pendingDeletion.Load()) + assert.False(t, hostinfo.isPendingDeletion()) assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0]) assert.Contains(t, nc.hostMap.Indexes, hostinfo.localIndexId) - assert.True(t, hostinfo.out.Load()) - assert.True(t, hostinfo.in.Load()) + assert.True(t, hostinfo.sentSinceCheck()) + assert.True(t, (hostinfo.state.Load()&stateIn != 0)) // Do a traffic check tick, should not be pending deletion but should not have any in/out packets recorded nc.doTrafficCheck(hostinfo.localIndexId, p, nb, out, time.Now()) - assert.False(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.False(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) // Do another traffic check tick, this host should be pending deletion now nc.Out(hostinfo) - assert.True(t, hostinfo.out.Load()) + assert.True(t, hostinfo.sentSinceCheck()) nc.doTrafficCheck(hostinfo.localIndexId, p, nb, out, time.Now()) - assert.True(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.True(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) assert.Contains(t, nc.hostMap.Indexes, hostinfo.localIndexId) assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0]) @@ -168,33 +168,33 @@ func Test_NewConnectionManagerTest2(t *testing.T) { // We saw traffic out to vpnIp nc.Out(hostinfo) nc.In(hostinfo) - assert.True(t, hostinfo.in.Load()) - assert.True(t, hostinfo.out.Load()) - assert.False(t, hostinfo.pendingDeletion.Load()) + assert.True(t, (hostinfo.state.Load()&stateIn != 0)) + assert.True(t, hostinfo.sentSinceCheck()) + assert.False(t, hostinfo.isPendingDeletion()) assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0]) assert.Contains(t, nc.hostMap.Indexes, hostinfo.localIndexId) // Do a traffic check tick, should not be pending deletion but should not have any in/out packets recorded nc.doTrafficCheck(hostinfo.localIndexId, p, nb, out, time.Now()) - assert.False(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.False(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) // Do another traffic check tick, this host should be pending deletion now nc.Out(hostinfo) nc.doTrafficCheck(hostinfo.localIndexId, p, nb, out, time.Now()) - assert.True(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.True(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) assert.Contains(t, nc.hostMap.Indexes, hostinfo.localIndexId) assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0]) // We saw traffic, should no longer be pending deletion nc.In(hostinfo) nc.doTrafficCheck(hostinfo.localIndexId, p, nb, out, time.Now()) - assert.False(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.False(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) assert.Contains(t, nc.hostMap.Indexes, hostinfo.localIndexId) assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0]) } @@ -253,31 +253,31 @@ func Test_NewConnectionManager_DisconnectInactive(t *testing.T) { // Do a traffic check tick, in and out should be cleared but should not be pending deletion nc.Out(hostinfo) nc.In(hostinfo) - assert.True(t, hostinfo.out.Load()) - assert.True(t, hostinfo.in.Load()) + assert.True(t, hostinfo.sentSinceCheck()) + assert.True(t, (hostinfo.state.Load()&stateIn != 0)) now := time.Now() decision, _, _ := nc.makeTrafficDecision(hostinfo.localIndexId, now) assert.Equal(t, tryRehandshake, decision) assert.Equal(t, now, hostinfo.lastUsed) - assert.False(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.False(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) decision, _, _ = nc.makeTrafficDecision(hostinfo.localIndexId, now.Add(time.Second*5)) assert.Equal(t, doNothing, decision) assert.Equal(t, now, hostinfo.lastUsed) - assert.False(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.False(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) // Do another traffic check tick, should still not be pending deletion decision, _, _ = nc.makeTrafficDecision(hostinfo.localIndexId, now.Add(time.Second*10)) assert.Equal(t, doNothing, decision) assert.Equal(t, now, hostinfo.lastUsed) - assert.False(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.False(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) assert.Contains(t, nc.hostMap.Indexes, hostinfo.localIndexId) assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0]) @@ -285,9 +285,9 @@ func Test_NewConnectionManager_DisconnectInactive(t *testing.T) { decision, _, _ = nc.makeTrafficDecision(hostinfo.localIndexId, now.Add(time.Minute*10)) assert.Equal(t, closeTunnel, decision) assert.Equal(t, now, hostinfo.lastUsed) - assert.False(t, hostinfo.pendingDeletion.Load()) - assert.False(t, hostinfo.out.Load()) - assert.False(t, hostinfo.in.Load()) + assert.False(t, hostinfo.isPendingDeletion()) + assert.False(t, hostinfo.sentSinceCheck()) + assert.False(t, (hostinfo.state.Load()&stateIn != 0)) assert.Contains(t, nc.hostMap.Indexes, hostinfo.localIndexId) assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0]) } diff --git a/control.go b/control.go index 7df5a09e..5d7b27c5 100644 --- a/control.go +++ b/control.go @@ -212,7 +212,7 @@ func (c *Control) RebindUDPServer() { c.f.lightHouse.SendUpdate() // Let the main interface know that we rebound so that underlying tunnels know to trigger punches from their remotes - c.f.rebindCount++ + c.f.rebindEpoch.Add(1) } // ListHostmapHosts returns details about the actual or pending (handshaking) hostmap by vpn ip diff --git a/control_tester.go b/control_tester.go index 546b9e87..c4f96f20 100644 --- a/control_tester.go +++ b/control_tester.go @@ -123,6 +123,16 @@ func (c *Control) SetLocalAddrsFn(fn func(*LocalAllowList) []netip.Addr) { c.f.lightHouse.localAddrsFn = fn } +// GetRebindEpochFor returns the rebind epoch a tunnel last sent under, so a test can tell whether a send +// consumed the epoch edge without having to infer it from lighthouse traffic. +func (c *Control) GetRebindEpochFor(vpnAddr netip.Addr) (uint32, bool) { + h := c.f.hostMap.QueryVpnAddr(vpnAddr) + if h == nil { + return 0, false + } + return h.state.Load() >> stateEpochShift, true +} + func (c *Control) KillPendingTunnel(vpnIp netip.Addr) bool { hostinfo := c.f.handshakeManager.QueryVpnAddr(vpnIp) if hostinfo == nil { diff --git a/e2e/rebind_test.go b/e2e/rebind_test.go index 2547f739..e415dd73 100644 --- a/e2e/rebind_test.go +++ b/e2e/rebind_test.go @@ -223,3 +223,57 @@ func TestRebindAdvertisesNewAddressAfterMove(t *testing.T) { lhControl.Stop() myControl.Stop() } + +// A relayed send records traffic but must not consume the rebind epoch. If it does, the next direct send to the +// relay host sees the epoch already current and never requeries, so the far side is never told to punch at our +// new address. This pins the SendVia call site, which the unit tests cannot reach. +func TestRebindRequeriesAfterRelayedSend(t *testing.T) { + t.Parallel() + ca, _, caKey, _ := cert_test.NewTestCaCert(cert.Version2, cert.Curve_CURVE25519, time.Now(), time.Now().Add(10*time.Minute), nil, nil, []string{}) + + // No lighthouse on purpose: it would hand out a direct address for them and nothing would relay. + myControl, myVpnIpNet, _, _ := newSimpleServer(cert.Version2, ca, caKey, "me", "10.128.0.1/24", m{"relay": m{"use_relays": true}}) + relayControl, relayVpnIpNet, relayUdpAddr, _ := newSimpleServer(cert.Version2, ca, caKey, "relay", "10.128.0.128/24", m{"relay": m{"am_relay": true}}) + theirControl, theirVpnIpNet, theirUdpAddr, _ := newSimpleServer(cert.Version2, ca, caKey, "them", "10.128.0.2/24", m{"relay": m{"use_relays": true}}) + + myControl.InjectLightHouseAddr(relayVpnIpNet[0].Addr(), relayUdpAddr) + myControl.InjectRelays(theirVpnIpNet[0].Addr(), []netip.Addr{relayVpnIpNet[0].Addr()}) + relayControl.InjectLightHouseAddr(theirVpnIpNet[0].Addr(), theirUdpAddr) + + r := router.NewR(t, myControl, relayControl, theirControl) + defer r.RenderFlow() + + myControl.Start() + relayControl.Start() + theirControl.Start() + + myControl.InjectTunPacket(BuildTunUDPPacket(theirVpnIpNet[0].Addr(), 80, myVpnIpNet[0].Addr(), 80, []byte("establish"))) + r.RouteForAllUntilTxTun(theirControl) + r.RouteFor(time.Millisecond * 500) + + hi := myControl.GetHostInfoByVpnAddr(theirVpnIpNet[0].Addr(), false) + require.NotNil(t, hi, "expected a tunnel to them") + require.NotEmpty(t, hi.CurrentRelaysToMe, "them must be reachable only via the relay for this test to mean anything") + // sendNoMetrics only reaches SendVia when there is no direct remote, so pin that too. Without this the test + // keeps passing while quietly sending direct and never exercising the relay path. + require.False(t, hi.CurrentRemote.IsValid(), "them must have no direct remote, otherwise SendVia is never called") + + before, ok := myControl.GetRebindEpochFor(relayVpnIpNet[0].Addr()) + require.True(t, ok, "expected a tunnel to the relay") + + myControl.RebindUDPServer() + + // Traffic to them goes through SendVia on the relay tunnel. That must record traffic without consuming the + // relay tunnel's own epoch edge, which belongs to the direct path. + myControl.InjectTunPacket(BuildTunUDPPacket(theirVpnIpNet[0].Addr(), 80, myVpnIpNet[0].Addr(), 80, []byte("relayed"))) + r.RouteForAllUntilTxTun(theirControl) + + after, ok := myControl.GetRebindEpochFor(relayVpnIpNet[0].Addr()) + require.True(t, ok) + assert.Equal(t, before, after, + "a relayed send consumed the relay tunnel's rebind epoch, so the next direct send will not requery") + + myControl.Stop() + relayControl.Stop() + theirControl.Stop() +} diff --git a/hostmap.go b/hostmap.go index 45515fc3..6466a9b9 100644 --- a/hostmap.go +++ b/hostmap.go @@ -262,11 +262,6 @@ type HostInfo struct { // This is used to limit lighthouse re-queries in chatty clients nextLHQuery atomic.Int64 - // lastRebindCount is the other side of Interface.rebindCount, if these values don't match then we need to ask LH - // for a punch from the remote end of this tunnel. The goal being to prime their conntrack for our traffic just like - // with a handshake - lastRebindCount int8 - // lastHandshakeTime records the time the remote side told us about at the stage when the handshake was completed locally // Stage 1 packet will contain it if I am a responder, stage 2 packet if I am an initiator // This is used to avoid an attack where a handshake packet is replayed after some time @@ -275,8 +270,8 @@ type HostInfo struct { lastRoam time.Time lastRoamRemote netip.AddrPort - //TODO: in, out, and others might benefit from being an atomic.Int32. We could collapse connectionManager pendingDeletion, relayUsed, and in/out into this 1 thing - in, out, pendingDeletion atomic.Bool + // Traffic bits, pendingDeletion, and the rebind epoch we last sent under + state atomic.Uint32 // lastUsed tracks the last time ConnectionManager checked the tunnel and it was in use. // This value will be behind against actual tunnel utilization in the hot path. @@ -658,7 +653,7 @@ func (hm *HostMap) unlockedAddHostInfo(hostinfo *HostInfo, f *Interface) { hm.Indexes[hostinfo.localIndexId] = hostinfo hm.RemoteIndexes[hostinfo.remoteIndexId] = hostinfo - hostinfo.out.Store(true) + hostinfo.markOut(f.rebindEpoch.Load()) if f.connectionManager != nil { // f.connectionManager is only nil in some unit tests f.connectionManager.trafficTimer.Add(hostinfo.localIndexId, f.connectionManager.checkInterval) } @@ -759,6 +754,68 @@ func (i *HostInfo) TryPromoteBest(preferredRanges []netip.Prefix, ifce *Interfac } } +// Bits within HostInfo.state, everything above stateEpochShift is the epoch +const ( + stateIn uint32 = 1 << iota + stateOut + statePendingDeletion + + stateFlags = stateIn | stateOut | statePendingDeletion + stateEpochShift = 3 +) + +// markIn records inbound traffic +func (i *HostInfo) markIn() { + if i.state.Load()&stateIn == 0 { + i.state.Or(stateIn) + } +} + +// markOut records a send and reports whether the epoch moved, meaning we want a punch from the far side +func (i *HostInfo) markOut(epoch uint32) bool { + e := epoch << stateEpochShift + for { + old := i.state.Load() + if old&stateOut != 0 && old&^stateFlags == e { + return false + } + + if i.state.CompareAndSwap(old, old&stateFlags|stateOut|e) { + return old&^stateFlags != e + } + } +} + +// markOutOnly records a send without consuming the rebind epoch, for paths that cannot act on a requery +func (i *HostInfo) markOutOnly() { + if i.state.Load()&stateOut == 0 { + i.state.Or(stateOut) + } +} + +// sentSinceCheck reports whether anything has been sent since the connection manager last looked +func (i *HostInfo) sentSinceCheck() bool { + return i.state.Load()&stateOut != 0 +} + +// takeTraffic clears both traffic bits, leaving the epoch alone, and reports what they were +func (i *HostInfo) takeTraffic() (in bool, out bool) { + old := i.state.And(^(stateIn | stateOut)) + return old&stateIn != 0, old&stateOut != 0 +} + +func (i *HostInfo) setPendingDeletion(v bool) { + if v { + i.state.Or(statePendingDeletion) + } else { + i.state.And(^statePendingDeletion) + } +} + +func (i *HostInfo) isPendingDeletion() bool { + return i.state.Load()&statePendingDeletion != 0 +} + func (i *HostInfo) GetCert() *cert.CachedCertificate { if i.ConnectionState != nil { return i.ConnectionState.peerCert diff --git a/hostmap_test.go b/hostmap_test.go index 9cfebe17..389af5ca 100644 --- a/hostmap_test.go +++ b/hostmap_test.go @@ -401,3 +401,43 @@ func TestHostMap_RelayState(t *testing.T) { assert.Equal(t, []netip.Addr{}, h1.relayState.relays) } + +func TestHostInfo_markOut(t *testing.T) { + h := &HostInfo{} + h.markOut(5) // stamped when the tunnel was added + + // A tunnel already on the current epoch has nothing to report, which is what keeps a fresh tunnel from + // requerying on its first packet + assert.False(t, h.markOut(5), "an unchanged epoch should not report a move") + assert.True(t, h.sentSinceCheck(), "the send is still recorded as traffic") + + // A rebind is observed exactly once, so we requery once per rebind + assert.True(t, h.markOut(6), "a bumped epoch should report a move") + assert.False(t, h.markOut(6), "the epoch move should only be reported once") + + // Traffic and pendingDeletion live in the same word and must survive an epoch change + h.setPendingDeletion(true) + h.markIn() + assert.True(t, h.markOut(7)) + assert.True(t, h.isPendingDeletion(), "pendingDeletion must survive an epoch change") + in, out := h.takeTraffic() + assert.True(t, in, "inbound traffic must survive an epoch change") + assert.True(t, out) + + // Clearing the traffic bits leaves the epoch alone, otherwise an idle tunnel would requery forever + assert.False(t, h.markOut(7), "takeTraffic must not disturb the epoch") +} + +// A relayed send records traffic but must leave the rebind epoch for the direct path to consume, otherwise +// relaying to a host swallows the requery that gets the far side punching at our new address. +func TestHostInfo_markOutOnly(t *testing.T) { + h := &HostInfo{} + h.markOut(5) + + h.markOutOnly() + assert.True(t, h.sentSinceCheck(), "a relayed send is still outbound traffic") + assert.False(t, h.markOut(5), "a relayed send must not disturb the epoch") + + assert.True(t, h.markOut(6), "a relayed send must not consume the epoch edge") + assert.False(t, h.markOut(6)) +} diff --git a/inside.go b/inside.go index a80b2e96..38361485 100644 --- a/inside.go +++ b/inside.go @@ -297,7 +297,7 @@ func (f *Interface) SendVia(via *HostInfo, c := via.ConnectionState.messageCounter.Add(1) out = header.Encode(out, header.Version, header.Message, header.MessageRelay, relay.RemoteIndex, c) - f.connectionManager.Out(via) + f.connectionManager.OutRelay(via) // Authenticate the header and payload, but do not encrypt for this message type. // The payload consists of the inner, unencrypted Nebula header, as well as the end-to-end encrypted payload. @@ -365,17 +365,11 @@ func (f *Interface) sendNoMetrics(t header.MessageType, st header.MessageSubType //l.WithField("trace", string(debug.Stack())).Error("out Header ", &Header{Version, t, st, 0, hostinfo.remoteIndexId, c}, p) out = header.Encode(out, header.Version, t, st, hostinfo.remoteIndexId, c) - f.connectionManager.Out(hostinfo) - - // Query our LH if we haven't since the last time we've been rebound, this will cause the remote to punch against - // all our addrs and enable a faster roaming. - if t != header.CloseTunnel && hostinfo.lastRebindCount != f.rebindCount { - //NOTE: there is an update hole if a tunnel isn't used and exactly 256 rebinds occur before the tunnel is - // finally used again. This tunnel would eventually be torn down and recreated if this action didn't help. + // We rebound since this tunnel last sent, ask the lighthouse to get the far side punching at us again + if f.connectionManager.Out(hostinfo) && t != header.CloseTunnel { f.lightHouse.QueryServer(hostinfo.vpnAddrs[0]) - hostinfo.lastRebindCount = f.rebindCount if f.l.Enabled(context.Background(), slog.LevelDebug) { - f.l.Debug("Lighthouse update triggered for punch due to rebind counter", + f.l.Debug("Lighthouse update triggered for punch due to rebind epoch", "vpnAddrs", hostinfo.vpnAddrs, ) } diff --git a/interface.go b/interface.go index c44f38b3..93141aea 100644 --- a/interface.go +++ b/interface.go @@ -82,8 +82,8 @@ type Interface struct { sendRecvErrorConfig recvErrorConfig acceptRecvErrorConfig recvErrorConfig - // rebindCount is used to decide if an active tunnel should trigger a punch notification through a lighthouse - rebindCount int8 + // Bumped on every udp rebind, tunnels compare it to decide they need a punch from the far side + rebindEpoch atomic.Uint32 version string conntrackCacheTimeout time.Duration