Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
fix: address tailnet coordinator review feedback
Contain cleanup panics at the connection goroutine boundary and stop compound requests when node update fan-out removes their peer. Clarify recovered mapping and nil handshake entry messages.

Generated by Coder Agents.
  • Loading branch information
dylanhuff-at-coder committed Sep 18, 2026
commit 24c87fad4d3e56b95397c0a31ea32b6c1168c98f
2 changes: 1 addition & 1 deletion enterprise/tailnet/connio_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ func TestConnIOHandleRequestRejectsBeforeMutation(t *testing.T) {
}},
ReadyForHandshake: []*proto.CoordinateRequest_ReadyForHandshake{nil},
})
require.EqualError(t, err, "ready_for_handshake entry is required")
require.EqualError(t, err, "ready_for_handshake entries must not be nil")
require.Contains(t, logbuf.String(), "invalid coordinate request")
select {
case binding := <-bindings:
Expand Down
2 changes: 1 addition & 1 deletion enterprise/tailnet/pgcoord.go
Original file line number Diff line number Diff line change
Expand Up @@ -745,7 +745,7 @@ func newMapper(c *connIO, logger slog.Logger, h *heartbeats) *mapper {
func (m *mapper) run() {
defer func() {
if recovered := recover(); recovered != nil {
m.logger.Error(m.ctx, "panic mapping peer responses (recovered)",
m.logger.Error(m.ctx, "panic processing peer mappings (recovered)",
slog.F("panic", recovered),
slog.F("stack", string(debug.Stack())),
)
Expand Down
14 changes: 14 additions & 0 deletions tailnet/coordinator.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"io"
"net/http"
"net/netip"
runtimedebug "runtime/debug"
"sync"
"time"

Expand Down Expand Up @@ -177,6 +178,15 @@ func (c *coordinator) Coordinate(
c.wg.Add(1)
go func() {
defer c.wg.Done()
// Peer cleanup can also panic, outside reqLoop's recovery boundary.
defer func() {
if recovered := recover(); recovered != nil {
logger.Error(ctx, "panic coordinating peer (recovered)",
slog.F("panic", recovered),
slog.F("stack", string(runtimedebug.Stack())),
)
}
}()
loopErr := p.reqLoop(ctx, logger, c.core.handleRequest)
closeErrStr := ""
if loopErr != nil {
Expand Down Expand Up @@ -357,6 +367,10 @@ func (c *core) nodeUpdateLocked(p *peer, node *proto.Node) (err error) {

p.node = node
c.updateTunnelPeersLocked(p.id, node, proto.CoordinateResponse_PeerUpdate_NODE, "node update")
// Fan-out can recursively remove this peer if its response buffer is full.
if c.peers[p.id] != p {
return ErrAlreadyRemoved
}
return nil
}

Expand Down
34 changes: 22 additions & 12 deletions tailnet/coordinator_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -163,8 +163,8 @@ func TestCoreDisconnectIgnoresReadyForHandshake(t *testing.T) {
// fails with ErrWouldBlock, which removes that peer too, and notifying back
// removes the first peer inside the nested call. The outer removePeerLocked
// must notice its peer is already gone instead of closing the channel again.
// lostPeer is the case that matters most: it runs on the Coordinate goroutine
// where there is no recover.
// Both request handling and lostPeer cleanup must finish without relying on
// panic recovery.
func TestCoreRemovePeerNestedRemoval(t *testing.T) {
t.Parallel()

Expand Down Expand Up @@ -248,19 +248,29 @@ func TestCoreRemovePeerNestedRemoval(t *testing.T) {
requireClosed(ctx, t, fp.bResps)
}

t.Run("UpdateSelf", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
fp := setup(t)
var err error
require.NotPanics(t, func() {
err = fp.core.handleRequest(ctx, fp.a, &proto.CoordinateRequest{
for _, name := range []string{"UpdateSelf", "UpdateSelfAddTunnel", "UpdateSelfReadyForHandshake"} {
t.Run(name, func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
fp := setup(t)
req := &proto.CoordinateRequest{
UpdateSelf: &proto.CoordinateRequest_UpdateSelf{Node: &proto.Node{PreferredDerp: 3}},
}
dstID := uuid.New()
switch name {
case "UpdateSelfAddTunnel":
req.AddTunnel = &proto.CoordinateRequest_Tunnel{Id: dstID[:]}
case "UpdateSelfReadyForHandshake":
req.ReadyForHandshake = []*proto.CoordinateRequest_ReadyForHandshake{{Id: dstID[:]}}
}
var err error
require.NotPanics(t, func() {
err = fp.core.handleRequest(ctx, fp.a, req)
})
require.NoError(t, err)
requireBothRemoved(ctx, t, fp)
})
require.NoError(t, err)
requireBothRemoved(ctx, t, fp)
})
}

t.Run("LostPeer", func(t *testing.T) {
t.Parallel()
Expand Down
71 changes: 71 additions & 0 deletions tailnet/peer_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"github.com/google/uuid"
"github.com/stretchr/testify/require"

"cdr.dev/slog/v3"
"cdr.dev/slog/v3/sloggers/slogtest"
"github.com/coder/coder/v2/tailnet/proto"
"github.com/coder/coder/v2/testutil"
Expand Down Expand Up @@ -46,3 +47,73 @@ func TestPeerReqLoopPanic(t *testing.T) {
core.mutex.RUnlock()
require.False(t, ok)
}

func TestPeerCleanupPanic(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitLong)
id := uuid.New()
sink := &peerCleanupPanicSink{
FakeSink: testutil.NewFakeSink(t),
peerID: id,
}
logger := slog.Make(sink).Leveled(slog.LevelDebug)
c := NewCoordinator(logger).(*coordinator)
t.Cleanup(func() { require.NoError(t, c.Close()) })

requests, responses := c.Coordinate(ctx, id, t.Name(), SingleTailnetCoordinateeAuth{})
close(requests)
done := make(chan struct{})
go func() {
c.wg.Wait()
close(done)
}()
testutil.TryReceive(ctx, t, done)

entries := sink.Entries(func(entry slog.SinkEntry) bool {
return entry.Level == slog.LevelError
})
require.Len(t, entries, 1)
require.Equal(t, "panic coordinating peer (recovered)", entries[0].Message)
fields := make(map[string]any)
for _, field := range entries[0].Fields {
fields[field.Name] = field.Value
}
require.Equal(t, id, fields["peer_id"])
require.Equal(t, "private cleanup panic", fields["panic"])
require.Contains(t, fields["stack"], "(*core).lostPeer")

// Recovery must leave the core unlocked and the coordinator usable.
healthyID := uuid.New()
healthyRequests, healthyResponses := c.Coordinate(ctx, healthyID, t.Name(), SingleTailnetCoordinateeAuth{})
close(healthyRequests)
require.Nil(t, testutil.TryReceive(ctx, t, healthyResponses))

// Shutdown must still remove the peer whose cleanup was interrupted.
closed := make(chan error, 1)
go func() { closed <- c.Close() }()
require.NoError(t, testutil.RequireReceive(ctx, t, closed))
response := testutil.RequireReceive(ctx, t, responses)
require.Equal(t, CloseErrCoordinatorClose, response.Error)
require.NotContains(t, response.Error, "private cleanup panic")
require.Nil(t, testutil.TryReceive(ctx, t, responses))
c.core.mutex.RLock()
_, present := c.core.peers[id]
c.core.mutex.RUnlock()
require.False(t, present)
}

type peerCleanupPanicSink struct {
*testutil.FakeSink
peerID uuid.UUID
}

func (s *peerCleanupPanicSink) LogEntry(ctx context.Context, entry slog.SinkEntry) {
if entry.Message == "lostPeer" {
for _, field := range entry.Fields {
if field.Name == "peer_id" && field.Value == s.peerID {
panic("private cleanup panic")
}
}
}
s.FakeSink.LogEntry(ctx, entry)
}
2 changes: 1 addition & 1 deletion tailnet/requests.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ func ValidateCoordinateRequest(req *proto.CoordinateRequest) error {
return err
}
if slices.Contains(req.ReadyForHandshake, nil) {
return xerrors.New("ready_for_handshake entry is required")
return xerrors.New("ready_for_handshake entries must not be nil")
}
return nil
}
Expand Down
4 changes: 2 additions & 2 deletions tailnet/requests_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,12 +79,12 @@ func TestValidateCoordinateRequest(t *testing.T) {
{
name: "NilReadyForHandshakeEntry",
req: &proto.CoordinateRequest{ReadyForHandshake: []*proto.CoordinateRequest_ReadyForHandshake{nil}},
err: "ready_for_handshake entry is required",
err: "ready_for_handshake entries must not be nil",
},
{
name: "NilReadyForHandshakeAfterValid",
req: &proto.CoordinateRequest{ReadyForHandshake: []*proto.CoordinateRequest_ReadyForHandshake{validRFH, nil}},
err: "ready_for_handshake entry is required",
err: "ready_for_handshake entries must not be nil",
},
} {
t.Run(tc.name, func(t *testing.T) {
Expand Down
4 changes: 2 additions & 2 deletions tailnet/test/requests.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ func InvalidCoordinateRequestTest(ctx context.Context, t *testing.T, coordinator
{
name: "NilReadyForHandshake",
req: &proto.CoordinateRequest{ReadyForHandshake: []*proto.CoordinateRequest_ReadyForHandshake{nil}},
err: "ready_for_handshake entry is required",
err: "ready_for_handshake entries must not be nil",
},
{
name: "NodeAndNilReadyForHandshake",
Expand All @@ -43,7 +43,7 @@ func InvalidCoordinateRequestTest(ctx context.Context, t *testing.T, coordinator
}},
ReadyForHandshake: []*proto.CoordinateRequest_ReadyForHandshake{nil},
},
err: "ready_for_handshake entry is required",
err: "ready_for_handshake entries must not be nil",
verifyNoMutation: true,
},
} {
Expand Down
Loading