mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
refactor: add postgres tailnet coordinator (#8044)
* postgres tailnet coordinator Signed-off-by: Spike Curtis <spike@coder.com> * Fix db migration; tests Signed-off-by: Spike Curtis <spike@coder.com> * Add fixture, regenerate Signed-off-by: Spike Curtis <spike@coder.com> * Fix fixtures Signed-off-by: Spike Curtis <spike@coder.com> * review comments, run clean gen Signed-off-by: Spike Curtis <spike@coder.com> * Rename waitForConn -> cleanupConn Signed-off-by: Spike Curtis <spike@coder.com> * code review updates Signed-off-by: Spike Curtis <spike@coder.com> * db migration order Signed-off-by: Spike Curtis <spike@coder.com> * fix log field name last_heartbeat Signed-off-by: Spike Curtis <spike@coder.com> * fix heartbeat_from log field Signed-off-by: Spike Curtis <spike@coder.com> * fix slog fields for linting Signed-off-by: Spike Curtis <spike@coder.com> --------- Signed-off-by: Spike Curtis <spike@coder.com>
This commit is contained in:
+12
-5
@@ -1,6 +1,7 @@
|
||||
package tailnet
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
@@ -174,11 +175,12 @@ func newCore(logger slog.Logger) *core {
|
||||
var ErrWouldBlock = xerrors.New("would block")
|
||||
|
||||
type TrackedConn struct {
|
||||
ctx context.Context
|
||||
cancel func()
|
||||
conn net.Conn
|
||||
updates chan []*Node
|
||||
logger slog.Logger
|
||||
ctx context.Context
|
||||
cancel func()
|
||||
conn net.Conn
|
||||
updates chan []*Node
|
||||
logger slog.Logger
|
||||
lastData []byte
|
||||
|
||||
// ID is an ephemeral UUID used to uniquely identify the owner of the
|
||||
// connection.
|
||||
@@ -224,6 +226,10 @@ func (t *TrackedConn) SendUpdates() {
|
||||
t.logger.Error(t.ctx, "unable to marshal nodes update", slog.Error(err), slog.F("nodes", nodes))
|
||||
return
|
||||
}
|
||||
if bytes.Equal(t.lastData, data) {
|
||||
t.logger.Debug(t.ctx, "skipping duplicate update", slog.F("nodes", nodes))
|
||||
continue
|
||||
}
|
||||
|
||||
// Set a deadline so that hung connections don't put back pressure on the system.
|
||||
// Node updates are tiny, so even the dinkiest connection can handle them if it's not hung.
|
||||
@@ -255,6 +261,7 @@ func (t *TrackedConn) SendUpdates() {
|
||||
_ = t.Close()
|
||||
return
|
||||
}
|
||||
t.lastData = data
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -96,7 +96,7 @@ func TestCoordinator(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
close(closeAgentChan)
|
||||
}()
|
||||
sendAgentNode(&tailnet.Node{})
|
||||
sendAgentNode(&tailnet.Node{PreferredDERP: 1})
|
||||
require.Eventually(t, func() bool {
|
||||
return coordinator.Node(agentID) != nil
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
@@ -122,7 +122,7 @@ func TestCoordinator(t *testing.T) {
|
||||
case <-ctx.Done():
|
||||
t.Fatal("timed out")
|
||||
}
|
||||
sendClientNode(&tailnet.Node{})
|
||||
sendClientNode(&tailnet.Node{PreferredDERP: 2})
|
||||
clientNodes := <-agentNodeChan
|
||||
require.Len(t, clientNodes, 1)
|
||||
|
||||
@@ -131,7 +131,7 @@ func TestCoordinator(t *testing.T) {
|
||||
time.Sleep(tailnet.WriteTimeout * 3 / 2)
|
||||
|
||||
// Ensure an update to the agent node reaches the client!
|
||||
sendAgentNode(&tailnet.Node{})
|
||||
sendAgentNode(&tailnet.Node{PreferredDERP: 3})
|
||||
select {
|
||||
case agentNodes := <-clientNodeChan:
|
||||
require.Len(t, agentNodes, 1)
|
||||
@@ -193,7 +193,7 @@ func TestCoordinator(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
close(closeAgentChan1)
|
||||
}()
|
||||
sendAgentNode1(&tailnet.Node{})
|
||||
sendAgentNode1(&tailnet.Node{PreferredDERP: 1})
|
||||
require.Eventually(t, func() bool {
|
||||
return coordinator.Node(agentID) != nil
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
@@ -215,12 +215,12 @@ func TestCoordinator(t *testing.T) {
|
||||
}()
|
||||
agentNodes := <-clientNodeChan
|
||||
require.Len(t, agentNodes, 1)
|
||||
sendClientNode(&tailnet.Node{})
|
||||
sendClientNode(&tailnet.Node{PreferredDERP: 2})
|
||||
clientNodes := <-agentNodeChan1
|
||||
require.Len(t, clientNodes, 1)
|
||||
|
||||
// Ensure an update to the agent node reaches the client!
|
||||
sendAgentNode1(&tailnet.Node{})
|
||||
sendAgentNode1(&tailnet.Node{PreferredDERP: 3})
|
||||
agentNodes = <-clientNodeChan
|
||||
require.Len(t, agentNodes, 1)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user