Skip to content

Commit 6638d0e

Browse files
authored
Allow bound tunnel fallback during endpoint health lag (#236)
* Allow bound tunnel fallback during endpoint health lag * Allow bound tunnel fallback during endpoint health lag * Fix bound tunnel fallback review issues
1 parent 72768d9 commit 6638d0e

2 files changed

Lines changed: 110 additions & 0 deletions

File tree

‎gravity/endpoint_independence_test.go‎

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1661,6 +1661,81 @@ func TestSelectStreamForPacket_MarksEndpointUnhealthy_TriggersPerEndpointReconne
16611661
}
16621662
}
16631663

1664+
// TestSelectStreamForPacket_UsesBoundTunnelWhenEndpointHealthIsStale verifies
1665+
// that response traffic can still use the endpoint it is already bound to when
1666+
// endpoint health has gone stale but the tunnel stream itself is still healthy.
1667+
// This covers the post-hello/post-tunnel reconnect window where
1668+
// refreshEndpointHealth() may mark every endpoint unhealthy before the data
1669+
// path has actually failed.
1670+
func TestSelectStreamForPacket_UsesBoundTunnelWhenEndpointHealthIsStale(t *testing.T) {
1671+
t.Parallel()
1672+
1673+
g := newWritePacketTestClient(t, 2)
1674+
1675+
g.endpoints[0].healthy.Store(false)
1676+
g.endpoints[1].healthy.Store(false)
1677+
g.streamManager.connectionHealth[0] = false
1678+
g.streamManager.connectionHealth[1] = false
1679+
g.streamManager.controlStreams[0] = &configurableMockStream{}
1680+
g.streamManager.controlStreams[1] = &configurableMockStream{}
1681+
1682+
boundStream := &hardeningMockTunnelStream{}
1683+
setupWritePacketStreams(g, []*StreamInfo{
1684+
{connIndex: 0, isHealthy: true, streamID: "ep0-t0", stream: boundStream, lastUsed: time.Now()},
1685+
{connIndex: 1, isHealthy: false, streamID: "ep1-t0", stream: &hardeningMockTunnelStream{}, lastUsed: time.Now()},
1686+
})
1687+
1688+
pkt := makeIPv6Packet()
1689+
preBindFlowToEndpoint(g, pkt, 0)
1690+
1691+
stream, err := g.selectStreamForPacket(pkt)
1692+
if err != nil {
1693+
t.Fatalf("expected bound healthy tunnel fallback, got: %v", err)
1694+
}
1695+
if stream.connIndex != 0 {
1696+
t.Fatalf("expected bound stream from endpoint 0, got connIndex=%d", stream.connIndex)
1697+
}
1698+
}
1699+
1700+
func TestSelectStreamForPacket_BoundTunnelFallbackRefreshesBindingTTL(t *testing.T) {
1701+
t.Parallel()
1702+
1703+
g := newWritePacketTestClient(t, 2)
1704+
1705+
g.endpoints[0].healthy.Store(false)
1706+
g.endpoints[1].healthy.Store(false)
1707+
g.streamManager.connectionHealth[0] = false
1708+
g.streamManager.connectionHealth[1] = false
1709+
g.streamManager.controlStreams[0] = &configurableMockStream{}
1710+
g.streamManager.controlStreams[1] = &configurableMockStream{}
1711+
1712+
boundStream := &hardeningMockTunnelStream{}
1713+
setupWritePacketStreams(g, []*StreamInfo{
1714+
{connIndex: 0, isHealthy: true, streamID: "ep0-t0", stream: boundStream, lastUsed: time.Now()},
1715+
})
1716+
1717+
pkt := makeIPv6Packet()
1718+
preBindFlowToEndpoint(g, pkt, 0)
1719+
1720+
key := ExtractFlowKey(pkt)
1721+
before := time.Now().Add(-2 * time.Second)
1722+
g.selector.mu.Lock()
1723+
g.selector.bindings[key].LastUsed = before
1724+
g.selector.mu.Unlock()
1725+
1726+
_, err := g.selectStreamForPacket(pkt)
1727+
if err != nil {
1728+
t.Fatalf("expected bound healthy tunnel fallback, got: %v", err)
1729+
}
1730+
1731+
g.selector.mu.RLock()
1732+
after := g.selector.bindings[key].LastUsed
1733+
g.selector.mu.RUnlock()
1734+
if !after.After(before) {
1735+
t.Fatalf("expected binding lastUsed to advance, before=%v after=%v", before, after)
1736+
}
1737+
}
1738+
16641739
// --- Category D: Safety Net Edge Cases ---
16651740

16661741
// TestTriggerAllEndpointReconnections_ClosingClient verifies that

‎gravity/grpc_client.go‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5313,6 +5313,10 @@ func (g *GravityClient) selectStreamForPacket(payload []byte) (*StreamInfo, erro
53135313
for attempt := 0; attempt < len(endpoints); attempt++ {
53145314
endpoint := selector.Select(payload, endpoints)
53155315
if endpoint == nil {
5316+
if fallbackStream, fallbackURL, ok := g.selectBoundTunnelFallback(payload, selector); ok {
5317+
g.logger.Debug("selectStream: using bound endpoint %s despite unhealthy endpoint state because a healthy tunnel stream still exists", fallbackURL)
5318+
return fallbackStream, nil
5319+
}
53165320
// Log endpoint health summary for debugging selector failures.
53175321
var healthSummary []string
53185322
for i, ep := range endpoints {
@@ -5337,6 +5341,10 @@ func (g *GravityClient) selectStreamForPacket(payload []byte) (*StreamInfo, erro
53375341
g.triggerEndpointReconnectByURL(endpoint.URL)
53385342
g.wakePeerDiscovery()
53395343
}
5344+
if fallbackStream, fallbackURL, ok := g.selectBoundTunnelFallback(payload, selector); ok {
5345+
g.logger.Debug("selectStream: using bound endpoint %s after selector exhausted healthy endpoint attempts because a healthy tunnel stream still exists", fallbackURL)
5346+
return fallbackStream, nil
5347+
}
53405348
return nil, fmt.Errorf("no healthy tunnel streams on any endpoint")
53415349
}
53425350

@@ -5370,6 +5378,33 @@ func (g *GravityClient) selectStreamForPacket(payload []byte) (*StreamInfo, erro
53705378
return stream, nil
53715379
}
53725380

5381+
func (g *GravityClient) selectBoundTunnelFallback(payload []byte, selector *EndpointSelector) (*StreamInfo, string, bool) {
5382+
if selector == nil {
5383+
return nil, "", false
5384+
}
5385+
5386+
key := ExtractFlowKey(payload)
5387+
now := time.Now()
5388+
5389+
selector.mu.RLock()
5390+
binding, ok := selector.bindings[key]
5391+
selector.mu.RUnlock()
5392+
if !ok || binding == nil || binding.Endpoint == nil || now.Sub(binding.LastUsed) >= selector.ttl {
5393+
return nil, "", false
5394+
}
5395+
5396+
stream, err := g.selectStreamForEndpoint(payload, binding.Endpoint.URL)
5397+
if err != nil {
5398+
return nil, binding.Endpoint.URL, false
5399+
}
5400+
selector.mu.Lock()
5401+
if current, ok := selector.bindings[key]; ok && current == binding {
5402+
current.LastUsed = now
5403+
}
5404+
selector.mu.Unlock()
5405+
return stream, binding.Endpoint.URL, true
5406+
}
5407+
53735408
func (g *GravityClient) selectStreamForEndpoint(payload []byte, endpointURL string) (*StreamInfo, error) {
53745409
g.mu.RLock()
53755410
candidateIndexes := append([]int(nil), g.endpointStreamIndices[endpointURL]...)

0 commit comments

Comments
 (0)