feat: batch connection logs to avoid DB lock contention (#23727)

- Running 30k connections was generating a ton of lock contention in the
DB
This commit is contained in:
Jon Ayers
2026-04-03 15:47:26 -05:00
committed by GitHub
parent 333503f74e
commit a1d51f0dab
21 changed files with 2168 additions and 426 deletions
+10 -2
View File
@@ -5,6 +5,7 @@ import (
"crypto/ed25519"
"crypto/tls"
"fmt"
"io"
"math"
"net/http"
"net/url"
@@ -144,10 +145,11 @@ func New(ctx context.Context, options *Options) (_ *API, err error) {
}
if options.ConnectionLogger == nil {
options.ConnectionLogger = connectionlog.NewConnectionLogger(
connectionlog.NewDBBackend(options.Database),
connLogger := connectionlog.New(
connectionlog.NewDBBatcher(ctx, options.Database, options.Logger),
connectionlog.NewSlogBackend(options.Logger),
)
options.ConnectionLogger = connLogger
}
meshTLSConfig, err := replicasync.CreateDERPMeshTLSConfig(options.AccessURL.Hostname(), options.TLSCertificates)
@@ -822,6 +824,12 @@ func (api *API) Close() error {
api.Options.CheckInactiveUsersCancelFunc()
}
// Close the connection logger to flush any remaining batched
// entries before shutting down the database connection.
if cl, ok := api.Options.ConnectionLogger.(io.Closer); ok {
_ = cl.Close()
}
return api.AGPL.Close()
}
+474 -16
View File
@@ -2,31 +2,70 @@ package connectionlog
import (
"context"
"io"
"sync"
"time"
"github.com/google/uuid"
"github.com/hashicorp/go-multierror"
"github.com/sqlc-dev/pqtype"
"cdr.dev/slog/v3"
agpl "github.com/coder/coder/v2/coderd/connectionlog"
"github.com/coder/coder/v2/coderd/database"
"github.com/coder/coder/v2/coderd/database/dbauthz"
auditbackends "github.com/coder/coder/v2/enterprise/audit/backends"
"github.com/coder/quartz"
)
const (
// defaultBatchSize is the maximum number of connection log entries
// to batch before forcing a flush.
defaultBatchSize = 1000
// defaultFlushInterval is how frequently to flush batched connection
// log entries to the database. Five seconds balances near-real-time
// audit visibility with write efficiency.
defaultFlushInterval = 5 * time.Second
// retryQueueSize is the capacity of the bounded retry channel.
// Failed batches beyond this limit are dropped.
retryQueueSize = 10
// shutdownWriteTimeout bounds how long a final write attempt
// can take during shutdown when the batcher context is already
// canceled.
shutdownWriteTimeout = 10 * time.Second
// maxRetries is the number of times to retry a failed batch
// write before dropping it and moving on.
maxRetries = 3
// retryInterval is the fixed delay between retry attempts.
retryInterval = time.Second
)
// Backend is a destination for connection log events. Backends that
// also implement io.Closer will be closed when the ConnectionLogger
// is closed.
type Backend interface {
Upsert(ctx context.Context, clog database.UpsertConnectionLogParams) error
}
func NewConnectionLogger(backends ...Backend) agpl.ConnectionLogger {
return &connectionLogger{
// ConnectionLogger fans out each connection log event to every
// registered backend.
type ConnectionLogger struct {
backends []Backend
}
// New creates a ConnectionLogger that dispatches to the given
// backends.
func New(backends ...Backend) *ConnectionLogger {
return &ConnectionLogger{
backends: backends,
}
}
type connectionLogger struct {
backends []Backend
}
func (c *connectionLogger) Upsert(ctx context.Context, clog database.UpsertConnectionLogParams) error {
func (c *ConnectionLogger) Upsert(ctx context.Context, clog database.UpsertConnectionLogParams) error {
var errs error
for _, backend := range c.backends {
err := backend.Upsert(ctx, clog)
@@ -37,24 +76,443 @@ func (c *connectionLogger) Upsert(ctx context.Context, clog database.UpsertConne
return errs
}
type dbBackend struct {
db database.Store
// Close closes all backends that implement io.Closer.
func (c *ConnectionLogger) Close() error {
var errs error
for _, backend := range c.backends {
if closer, ok := backend.(io.Closer); ok {
if err := closer.Close(); err != nil {
errs = multierror.Append(errs, err)
}
}
}
return errs
}
func NewDBBackend(db database.Store) Backend {
return &dbBackend{db: db}
// DBBatcherOption is a functional option for configuring a DBBatcher.
type DBBatcherOption func(b *DBBatcher)
// WithBatchSize sets the maximum number of entries to accumulate
// before forcing a flush.
func WithBatchSize(size int) DBBatcherOption {
return func(b *DBBatcher) {
b.maxBatchSize = size
}
}
func (b *dbBackend) Upsert(ctx context.Context, clog database.UpsertConnectionLogParams) error {
//nolint:gocritic // This is the Connection Logger
_, err := b.db.UpsertConnectionLog(dbauthz.AsConnectionLogger(ctx), clog)
return err
// WithFlushInterval sets how frequently the batcher flushes to the
// database.
func WithFlushInterval(d time.Duration) DBBatcherOption {
return func(b *DBBatcher) {
b.interval = d
}
}
// WithClock sets the clock, useful for testing.
func WithClock(clock quartz.Clock) DBBatcherOption {
return func(b *DBBatcher) {
b.clock = clock
}
}
// DBBatcher batches connection log upserts and periodically flushes
// them to the database to reduce per-event write pressure.
type DBBatcher struct {
store database.Store
log slog.Logger
itemCh chan database.UpsertConnectionLogParams
// dedupedBatch holds entries keyed by connection ID so that
// PostgreSQL never sees the same row twice in one INSERT …
// ON CONFLICT DO UPDATE. Connection IDs are globally unique
// (each new session gets a fresh UUID). Entries with a NULL
// connection_id (web events) go into nullConnIDBatch instead
// because NULL != NULL in SQL unique constraints.
dedupedBatch map[uuid.UUID]batchEntry
nullConnIDBatch []batchEntry
maxBatchSize int
// retryCh is a bounded channel of failed batches awaiting
// retry. A single retry worker goroutine processes this
// channel, retrying each batch up to maxRetries times before
// dropping it. If the channel is full, new failures are
// dropped immediately.
retryCh chan database.BatchUpsertConnectionLogsParams
clock quartz.Clock
timer *quartz.Timer
interval time.Duration
ctx context.Context
cancel context.CancelFunc
wg sync.WaitGroup
}
// NewDBBatcher creates a DBBatcher that batches writes to the database
// and starts its background processing loop. Close must be called to
// flush remaining entries on shutdown.
func NewDBBatcher(ctx context.Context, store database.Store, log slog.Logger, opts ...DBBatcherOption) *DBBatcher {
b := &DBBatcher{
store: store,
log: log,
clock: quartz.NewReal(),
}
for _, opt := range opts {
opt(b)
}
if b.interval == 0 {
b.interval = defaultFlushInterval
}
if b.maxBatchSize == 0 {
b.maxBatchSize = defaultBatchSize
}
b.timer = b.clock.NewTimer(b.interval)
b.itemCh = make(chan database.UpsertConnectionLogParams, b.maxBatchSize)
b.dedupedBatch = make(map[uuid.UUID]batchEntry, b.maxBatchSize)
b.retryCh = make(chan database.BatchUpsertConnectionLogsParams, retryQueueSize)
b.ctx, b.cancel = context.WithCancel(ctx)
b.wg.Add(2)
go func() {
defer b.wg.Done()
b.run(b.ctx)
}()
go func() {
defer b.wg.Done()
b.retryLoop()
}()
return b
}
// Upsert enqueues a connection log entry for batched writing. It
// blocks if the internal buffer is full, ensuring no logs are dropped.
// It returns an error if the batcher or caller context is canceled.
func (b *DBBatcher) Upsert(ctx context.Context, clog database.UpsertConnectionLogParams) error {
if b.ctx.Err() != nil {
return b.ctx.Err()
}
select {
case b.itemCh <- clog:
return nil
case <-b.ctx.Done():
return b.ctx.Err()
case <-ctx.Done():
return ctx.Err()
}
}
// Close cancels the batcher context, waits for the run loop and
// retry worker to exit.
func (b *DBBatcher) Close() error {
b.cancel()
if b.timer != nil {
b.timer.Stop()
}
b.wg.Wait()
return nil
}
// addToBatch inserts an item into the batch, deduplicating by conflict
// key on the fly. For entries with the same key, disconnect events are
// preferred over connect events, and later events are preferred over
// earlier ones.
//
// This is safe because each new connection gets a fresh UUID (see
// agent/agent.go and agent/agentssh), so the only duplicate for the
// same (connection_id, workspace_id, agent_name) is a connect/disconnect
// pair for the same session. A "reconnect" always uses a new ID.
func (b *DBBatcher) addToBatch(item database.UpsertConnectionLogParams) {
entry := batchEntry{
UpsertConnectionLogParams: item,
}
if item.ConnectionStatus == database.ConnectionStatusDisconnected {
// For standalone disconnect events, use the disconnect
// time as both connect and disconnect time. This matches
// the single-row UpsertConnectionLog behavior which uses
// @time for connect_time regardless of status. The SQL
// LEAST logic will correct connect_time if the real
// connect event arrives in a later batch.
entry.connectTime = item.Time
entry.disconnectTime = item.Time
} else {
entry.connectTime = item.Time
}
if !item.ConnectionID.Valid {
b.nullConnIDBatch = append(b.nullConnIDBatch, entry)
return
}
connID := item.ConnectionID.UUID
existing, ok := b.dedupedBatch[connID]
if !ok {
b.dedupedBatch[connID] = entry
return
}
// When merging entries for the same connection, always preserve
// the earliest non-zero connect_time and latest disconnect_time
// so the row records the full session span.
if !existing.connectTime.IsZero() && existing.connectTime.Before(entry.connectTime) {
entry.connectTime = existing.connectTime
}
if existing.disconnectTime.After(entry.disconnectTime) {
entry.disconnectTime = existing.disconnectTime
}
// Prefer disconnect over connect (superset of info).
// If same status, prefer the later event.
if item.ConnectionStatus == database.ConnectionStatusDisconnected &&
existing.ConnectionStatus != database.ConnectionStatusDisconnected {
b.dedupedBatch[connID] = entry
} else if item.Time.After(existing.Time) {
b.dedupedBatch[connID] = entry
}
}
// batchLen returns the total number of entries currently buffered.
func (b *DBBatcher) batchLen() int {
return len(b.dedupedBatch) + len(b.nullConnIDBatch)
}
func (b *DBBatcher) run(ctx context.Context) {
//nolint:gocritic // System-level batch operation for connection logs.
authCtx := dbauthz.AsConnectionLogger(ctx)
for ctx.Err() == nil {
select {
case item := <-b.itemCh:
b.addToBatch(item)
if b.batchLen() >= b.maxBatchSize {
b.flush(authCtx)
b.timer.Reset(b.interval, "connectionLogBatcher", "capacityFlush")
}
case <-b.timer.C:
b.flush(authCtx)
b.timer.Reset(b.interval, "connectionLogBatcher", "scheduledFlush")
case <-ctx.Done():
}
}
b.log.Debug(ctx, "context done, flushing before exit")
// Drain any remaining items from the channel.
for {
select {
case item := <-b.itemCh:
b.addToBatch(item)
default:
if b.batchLen() > 0 {
b.shutdownBatch(b.buildParams())
}
// Signal the retry worker to skip delays and close
// the channel so it exits after processing any
// remaining items.
// Mark the batcher as closed so that any subsequent
// Upsert calls fail immediately instead of sending
// into itemCh after the run loop has exited.
close(b.retryCh)
return
}
}
}
// batchEntry wraps a connection log event with explicit connect and
// disconnect times. When a connect and disconnect for the same session
// are merged into one entry, connectTime preserves the original
// session start while disconnectTime records when it ended.
type batchEntry struct {
database.UpsertConnectionLogParams
connectTime time.Time
disconnectTime time.Time
}
// flush builds the batch params, clears the in-memory batch, and
// writes to the database. On failure, the batch is queued for retry
// by the single retry worker goroutine. If the retry queue is full,
// the batch is dropped.
func (b *DBBatcher) flush(ctx context.Context) {
count := b.batchLen()
if count == 0 {
return
}
params := b.buildParams()
// Clear the batch before writing so the run loop can start
// accumulating new entries.
b.dedupedBatch = make(map[uuid.UUID]batchEntry, b.maxBatchSize)
b.nullConnIDBatch = nil
// Use the batcher's context for normal operation so Close()
// can cancel hung writes. During shutdown (ctx already canceled),
// fall back to a bounded timeout.
writeCtx := b.ctx
if writeCtx.Err() != nil {
var cancel context.CancelFunc
writeCtx, cancel = context.WithTimeout(context.Background(), shutdownWriteTimeout)
defer cancel()
}
//nolint:gocritic // System-level batch operation for connection logs.
err := b.store.BatchUpsertConnectionLogs(dbauthz.AsConnectionLogger(writeCtx), params)
if err == nil {
return
}
b.log.Error(ctx, "batch upsert failed, queueing for retry",
slog.Error(err), slog.F("count", count))
// Don't retry on shutdown.
if ctx.Err() != nil {
return
}
select {
case b.retryCh <- params:
default:
b.log.Error(ctx, "retry queue full, dropping batch",
slog.F("dropped", count))
}
}
func (b *DBBatcher) buildParams() database.BatchUpsertConnectionLogsParams {
count := b.batchLen()
var (
ids = make([]uuid.UUID, 0, count)
connectTime = make([]time.Time, 0, count)
organizationID = make([]uuid.UUID, 0, count)
workspaceOwnerID = make([]uuid.UUID, 0, count)
workspaceID = make([]uuid.UUID, 0, count)
workspaceName = make([]string, 0, count)
agentName = make([]string, 0, count)
connType = make([]database.ConnectionType, 0, count)
code = make([]int32, 0, count)
codeValid = make([]bool, 0, count)
ip = make([]pqtype.Inet, 0, count)
userAgent = make([]string, 0, count)
userID = make([]uuid.UUID, 0, count)
slugOrPort = make([]string, 0, count)
connectionID = make([]uuid.UUID, 0, count)
disconnectReason = make([]string, 0, count)
disconnectTime = make([]time.Time, 0, count)
)
appendEntry := func(e batchEntry) {
ids = append(ids, e.ID)
connectTime = append(connectTime, e.connectTime)
organizationID = append(organizationID, e.OrganizationID)
workspaceOwnerID = append(workspaceOwnerID, e.WorkspaceOwnerID)
workspaceID = append(workspaceID, e.WorkspaceID)
workspaceName = append(workspaceName, e.WorkspaceName)
agentName = append(agentName, e.AgentName)
connType = append(connType, e.Type)
code = append(code, e.Code.Int32)
codeValid = append(codeValid, e.Code.Valid)
ip = append(ip, e.IP)
userAgent = append(userAgent, e.UserAgent.String)
userID = append(userID, e.UserID.UUID)
slugOrPort = append(slugOrPort, e.SlugOrPort.String)
connectionID = append(connectionID, e.ConnectionID.UUID)
disconnectReason = append(disconnectReason, e.DisconnectReason.String)
disconnectTime = append(disconnectTime, e.disconnectTime)
}
for _, entry := range b.dedupedBatch {
appendEntry(entry)
}
for _, entry := range b.nullConnIDBatch {
appendEntry(entry)
}
return database.BatchUpsertConnectionLogsParams{
ID: ids,
ConnectTime: connectTime,
OrganizationID: organizationID,
WorkspaceOwnerID: workspaceOwnerID,
WorkspaceID: workspaceID,
WorkspaceName: workspaceName,
AgentName: agentName,
Type: connType,
Code: code,
CodeValid: codeValid,
Ip: ip,
UserAgent: userAgent,
UserID: userID,
SlugOrPort: slugOrPort,
ConnectionID: connectionID,
DisconnectReason: disconnectReason,
DisconnectTime: disconnectTime,
}
}
// retryLoop is a single background goroutine that processes failed
// batches from retryCh. Each batch is retried up to maxRetries times
// with a fixed delay between attempts. When draining is set (shutdown),
// batches get a single immediate write attempt instead. The loop exits
// when retryCh is closed by the run goroutine.
func (b *DBBatcher) retryLoop() {
for params := range b.retryCh {
b.retryBatch(params)
}
}
// retryBatch retries writing a batch up to maxRetries times with a
// fixed delay between attempts. If the batcher context is canceled
// during a wait, one final attempt is made before returning.
func (b *DBBatcher) retryBatch(params database.BatchUpsertConnectionLogsParams) {
count := len(params.ID)
for attempt := range maxRetries {
t := time.NewTimer(retryInterval)
select {
case <-b.ctx.Done():
b.shutdownBatch(params)
return
case <-t.C:
}
//nolint:gocritic // System-level batch operation for connection logs.
err := b.store.BatchUpsertConnectionLogs(dbauthz.AsConnectionLogger(b.ctx), params)
if err == nil {
return
}
b.log.Warn(b.ctx, "batch retry failed",
slog.Error(err),
slog.F("count", count),
slog.F("attempt", attempt+1),
slog.F("max_attempts", maxRetries),
)
}
b.log.Error(b.ctx, "batch retries exhausted, dropping batch",
slog.F("dropped", count))
}
// shutdownBatch makes a single write attempt during shutdown with a
// bounded timeout so it can't hang indefinitely.
func (b *DBBatcher) shutdownBatch(params database.BatchUpsertConnectionLogsParams) {
ctx, cancel := context.WithTimeout(context.Background(), shutdownWriteTimeout)
defer cancel()
//nolint:gocritic // System-level batch operation for connection logs.
err := b.store.BatchUpsertConnectionLogs(dbauthz.AsConnectionLogger(ctx), params)
if err != nil {
b.log.Error(b.ctx, "batch write failed on shutdown, dropping batch",
slog.Error(err), slog.F("dropped", len(params.ID)))
}
}
type connectionSlogBackend struct {
exporter *auditbackends.SlogExporter
}
// NewSlogBackend returns a Backend that logs connection events via
// the structured logger.
func NewSlogBackend(logger slog.Logger) Backend {
return &connectionSlogBackend{
exporter: auditbackends.NewSlogExporter(logger),
@@ -0,0 +1,529 @@
package connectionlog
import (
"context"
"database/sql"
"sync"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/require"
"go.uber.org/mock/gomock"
"golang.org/x/xerrors"
"cdr.dev/slog/v3"
"cdr.dev/slog/v3/sloggers/slogtest"
"github.com/coder/coder/v2/coderd/database"
"github.com/coder/coder/v2/coderd/database/dbmock"
"github.com/coder/coder/v2/testutil"
"github.com/coder/quartz"
)
func Test_addToBatch(t *testing.T) {
t.Parallel()
t.Run("ConnectThenDisconnect", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
wsID := uuid.New()
connID := uuid.New()
connect := fakeConnectEvent(wsID, "agent1", connID)
disconnect := fakeDisconnectEvent(wsID, "agent1", connID)
b.addToBatch(connect)
b.addToBatch(disconnect)
require.Equal(t, 1, b.batchLen())
key := connID
got := b.dedupedBatch[key]
require.Equal(t, disconnect.ID, got.ID)
require.Equal(t, database.ConnectionStatusDisconnected, got.ConnectionStatus)
// The connect_time should be preserved from the original
// connect event, not overwritten by the disconnect's
// timestamp.
require.Equal(t, connect.Time, got.connectTime)
require.Equal(t, disconnect.Time, got.disconnectTime)
})
t.Run("DisconnectThenLaterConnect", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
wsID := uuid.New()
connID := uuid.New()
disconnect := fakeDisconnectEvent(wsID, "agent1", connID)
connect := fakeConnectEvent(wsID, "agent1", connID)
connect.Time = disconnect.Time.Add(time.Second)
b.addToBatch(disconnect)
b.addToBatch(connect)
require.Equal(t, 1, b.batchLen())
key := connID
// The later event wins when the incoming item is not a
// disconnect. In practice, this case doesn't occur because
// connection IDs are never reused.
got := b.dedupedBatch[key]
require.Equal(t, connect.ID, got.ID)
// The disconnect's time should be preserved even though
// the connect event replaced it.
require.Equal(t, disconnect.Time, got.disconnectTime)
})
t.Run("DisconnectThenEarlierConnect", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
wsID := uuid.New()
connID := uuid.New()
disconnect := fakeDisconnectEvent(wsID, "agent1", connID)
connect := fakeConnectEvent(wsID, "agent1", connID)
connect.Time = disconnect.Time.Add(-time.Second)
b.addToBatch(disconnect)
b.addToBatch(connect)
require.Equal(t, 1, b.batchLen())
key := connID
require.Equal(t, disconnect.ID, b.dedupedBatch[key].ID)
})
t.Run("SameStatusKeepsLater", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
wsID := uuid.New()
connID := uuid.New()
early := fakeConnectEvent(wsID, "agent1", connID)
early.Time = time.Now()
late := fakeConnectEvent(wsID, "agent1", connID)
late.Time = early.Time.Add(time.Second)
b.addToBatch(early)
b.addToBatch(late)
require.Equal(t, 1, b.batchLen())
key := connID
require.Equal(t, late.ID, b.dedupedBatch[key].ID)
})
t.Run("NullConnIDsNeverDedup", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
evt1 := fakeNullConnIDEvent()
evt2 := fakeNullConnIDEvent()
evt2.WorkspaceID = evt1.WorkspaceID
evt2.AgentName = evt1.AgentName
b.addToBatch(evt1)
b.addToBatch(evt2)
require.Equal(t, 2, b.batchLen())
require.Len(t, b.nullConnIDBatch, 2)
require.Empty(t, b.dedupedBatch)
})
t.Run("MixedNullAndNonNull", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
wsID := uuid.New()
regular := fakeConnectEvent(wsID, "agent1", uuid.New())
nullEvt := fakeNullConnIDEvent()
nullEvt.WorkspaceID = wsID
nullEvt.AgentName = "agent1"
b.addToBatch(regular)
b.addToBatch(nullEvt)
require.Equal(t, 2, b.batchLen())
require.Len(t, b.dedupedBatch, 1)
require.Len(t, b.nullConnIDBatch, 1)
})
t.Run("StandaloneDisconnectUsesTimeAsConnectTime", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
connID := uuid.New()
disconnect := fakeDisconnectEvent(uuid.New(), "agent1", connID)
b.addToBatch(disconnect)
got := b.dedupedBatch[connID]
// A standalone disconnect must not leave connectTime as
// zero — that would insert a year-0001 connect_time in
// the DB. It should use the disconnect's own timestamp,
// matching the single-row UpsertConnectionLog behavior.
require.False(t, got.connectTime.IsZero(),
"standalone disconnect must have non-zero connectTime")
require.Equal(t, disconnect.Time, got.connectTime)
require.Equal(t, disconnect.Time, got.disconnectTime)
})
t.Run("DuplicateDisconnectsPreserveConnectTime", func(t *testing.T) {
t.Parallel()
b := &DBBatcher{
maxBatchSize: 100,
dedupedBatch: make(map[uuid.UUID]batchEntry),
}
wsID := uuid.New()
connID := uuid.New()
connect := fakeConnectEvent(wsID, "agent1", connID)
disconnect1 := fakeDisconnectEvent(wsID, "agent1", connID)
disconnect2 := fakeDisconnectEvent(wsID, "agent1", connID)
disconnect2.Time = disconnect1.Time.Add(time.Second)
b.addToBatch(connect)
b.addToBatch(disconnect1)
b.addToBatch(disconnect2)
require.Equal(t, 1, b.batchLen())
got := b.dedupedBatch[connID]
// The second disconnect should win (later event) but the
// original connect_time from the connect event must be
// preserved, not regressed to the disconnect's timestamp.
require.Equal(t, disconnect2.ID, got.ID)
require.Equal(t, connect.Time, got.connectTime,
"connect_time must not regress to disconnect timestamp")
require.Equal(t, disconnect2.Time, got.disconnectTime)
})
}
func Test_batcherFlush(t *testing.T) {
t.Parallel()
t.Run("DeduplicatesConnectDisconnect", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
ctrl := gomock.NewController(t)
store := dbmock.NewMockStore(ctrl)
clock := quartz.NewMock(t)
b := NewDBBatcher(ctx, store, log, WithClock(clock), WithBatchSize(100))
wsID := uuid.New()
connID := uuid.New()
connect := fakeConnectEvent(wsID, "agent1", connID)
disconnect := fakeDisconnectEvent(wsID, "agent1", connID)
// Expect a single batch with only the disconnect event.
store.EXPECT().
BatchUpsertConnectionLogs(gomock.Any(), batchParamsMatcher{
expectedCount: 1,
mustContainIDs: []uuid.UUID{disconnect.ID},
mustNotContainIDs: []uuid.UUID{connect.ID},
}).
Return(nil).
Times(1)
require.NoError(t, b.Upsert(ctx, connect))
require.NoError(t, b.Upsert(ctx, disconnect))
require.NoError(t, b.Close())
})
t.Run("DoesNotDeduplicateNullConnIDs", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
ctrl := gomock.NewController(t)
store := dbmock.NewMockStore(ctrl)
clock := quartz.NewMock(t)
b := NewDBBatcher(ctx, store, log, WithClock(clock), WithBatchSize(100))
evt1 := fakeNullConnIDEvent()
evt2 := fakeNullConnIDEvent()
evt2.WorkspaceID = evt1.WorkspaceID
evt2.AgentName = evt1.AgentName
store.EXPECT().
BatchUpsertConnectionLogs(gomock.Any(), batchParamsMatcher{
expectedCount: 2,
mustContainIDs: []uuid.UUID{evt1.ID, evt2.ID},
}).
Return(nil).
Times(1)
require.NoError(t, b.Upsert(ctx, evt1))
require.NoError(t, b.Upsert(ctx, evt2))
require.NoError(t, b.Close())
})
t.Run("DoesNotDeduplicateDifferentConnectionIDs", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
ctrl := gomock.NewController(t)
store := dbmock.NewMockStore(ctrl)
clock := quartz.NewMock(t)
b := NewDBBatcher(ctx, store, log, WithClock(clock), WithBatchSize(100))
wsID := uuid.New()
evt1 := fakeConnectEvent(wsID, "agent1", uuid.New())
evt2 := fakeConnectEvent(wsID, "agent1", uuid.New())
store.EXPECT().
BatchUpsertConnectionLogs(gomock.Any(), batchParamsMatcher{
expectedCount: 2,
mustContainIDs: []uuid.UUID{evt1.ID, evt2.ID},
}).
Return(nil).
Times(1)
require.NoError(t, b.Upsert(ctx, evt1))
require.NoError(t, b.Upsert(ctx, evt2))
require.NoError(t, b.Close())
})
t.Run("CloseFlushesMultipleEvents", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
ctrl := gomock.NewController(t)
store := dbmock.NewMockStore(ctrl)
clock := quartz.NewMock(t)
b := NewDBBatcher(ctx, store, log, WithClock(clock), WithBatchSize(100))
evt1 := fakeConnectEvent(uuid.New(), "agent1", uuid.New())
evt2 := fakeConnectEvent(uuid.New(), "agent2", uuid.New())
store.EXPECT().
BatchUpsertConnectionLogs(gomock.Any(), batchParamsMatcher{
expectedCount: 2,
mustContainIDs: []uuid.UUID{evt1.ID, evt2.ID},
}).
Return(nil).
Times(1)
require.NoError(t, b.Upsert(ctx, evt1))
require.NoError(t, b.Upsert(ctx, evt2))
require.NoError(t, b.Close())
})
t.Run("RetriesOnTransientFailure", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
ctrl := gomock.NewController(t)
store := dbmock.NewMockStore(ctrl)
clock := quartz.NewMock(t)
scheduledTrap := clock.Trap().TimerReset("connectionLogBatcher", "scheduledFlush")
defer scheduledTrap.Close()
b := NewDBBatcher(ctx, store, log, WithClock(clock), WithBatchSize(100))
evt := fakeConnectEvent(uuid.New(), "agent1", uuid.New())
// First call (synchronous in flush) fails, then the
// retry worker retries after the backoff and succeeds.
gomock.InOrder(
store.EXPECT().
BatchUpsertConnectionLogs(gomock.Any(), gomock.Any()).
Return(xerrors.New("transient error")).
Times(1),
store.EXPECT().
BatchUpsertConnectionLogs(gomock.Any(), batchParamsMatcher{
expectedCount: 1,
mustContainIDs: []uuid.UUID{evt.ID},
}).
Return(nil).
Times(1),
)
require.NoError(t, b.Upsert(ctx, evt))
// Trigger a scheduled flush while the batcher is still
// running. The synchronous write fails and queues to
// retryCh. The retry worker picks it up after a real-
// time 1s delay and succeeds.
clock.Advance(defaultFlushInterval).MustWait(ctx)
scheduledTrap.MustWait(ctx).MustRelease(ctx)
// Wait for the retry to complete (real-time 1s delay).
require.NoError(t, b.Close())
})
t.Run("ShutdownDrainsRetryQueue", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
ctrl := gomock.NewController(t)
store := dbmock.NewMockStore(ctrl)
clock := quartz.NewMock(t)
scheduledTrap := clock.Trap().TimerReset("connectionLogBatcher", "scheduledFlush")
defer scheduledTrap.Close()
b := NewDBBatcher(ctx, store, log, WithClock(clock), WithBatchSize(100))
evt := fakeConnectEvent(uuid.New(), "agent1", uuid.New())
// Track all successfully written IDs.
var writtenIDs []uuid.UUID
var mu sync.Mutex
firstCall := true
store.EXPECT().
BatchUpsertConnectionLogs(gomock.Any(), gomock.Any()).
DoAndReturn(func(_ context.Context, p database.BatchUpsertConnectionLogsParams) error {
mu.Lock()
defer mu.Unlock()
// First call (synchronous flush) fails, queueing
// the batch for retry.
if firstCall {
firstCall = false
return xerrors.New("transient error")
}
// Drain/retry attempts succeed.
writtenIDs = append(writtenIDs, p.ID...)
return nil
}).
AnyTimes()
// Send event and trigger flush — fails, queues.
require.NoError(t, b.Upsert(ctx, evt))
clock.Advance(defaultFlushInterval).MustWait(ctx)
scheduledTrap.MustWait(ctx).MustRelease(ctx)
// Close triggers shutdown. The retry worker drains
// retryCh and writes the batch via writeBatch.
require.NoError(t, b.Close())
mu.Lock()
defer mu.Unlock()
require.Contains(t, writtenIDs, evt.ID,
"event should be written during shutdown drain")
})
}
// batchParamsMatcher validates BatchUpsertConnectionLogsParams by
// checking count and specific IDs.
type batchParamsMatcher struct {
expectedCount int
mustContainIDs []uuid.UUID
mustNotContainIDs []uuid.UUID
}
func (m batchParamsMatcher) Matches(x interface{}) bool {
params, ok := x.(database.BatchUpsertConnectionLogsParams)
if !ok {
return false
}
if m.expectedCount > 0 && len(params.ID) != m.expectedCount {
return false
}
idSet := make(map[uuid.UUID]struct{}, len(params.ID))
for _, id := range params.ID {
idSet[id] = struct{}{}
}
for _, id := range m.mustContainIDs {
if _, ok := idSet[id]; !ok {
return false
}
}
for _, id := range m.mustNotContainIDs {
if _, ok := idSet[id]; ok {
return false
}
}
return true
}
func (batchParamsMatcher) String() string {
return "batch upsert params matcher"
}
func fakeConnectEvent(workspaceID uuid.UUID, agentName string, connectionID uuid.UUID) database.UpsertConnectionLogParams {
return database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: time.Now(),
OrganizationID: uuid.New(),
WorkspaceOwnerID: uuid.New(),
WorkspaceID: workspaceID,
WorkspaceName: "test-workspace",
AgentName: agentName,
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: connectionID, Valid: true},
ConnectionStatus: database.ConnectionStatusConnected,
}
}
func fakeDisconnectEvent(workspaceID uuid.UUID, agentName string, connectionID uuid.UUID) database.UpsertConnectionLogParams {
return database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: time.Now().Add(time.Second),
OrganizationID: uuid.New(),
WorkspaceOwnerID: uuid.New(),
WorkspaceID: workspaceID,
WorkspaceName: "test-workspace",
AgentName: agentName,
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: connectionID, Valid: true},
ConnectionStatus: database.ConnectionStatusDisconnected,
Code: sql.NullInt32{Int32: 0, Valid: true},
DisconnectReason: sql.NullString{String: "normal", Valid: true},
}
}
func fakeNullConnIDEvent() database.UpsertConnectionLogParams {
return database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: time.Now(),
OrganizationID: uuid.New(),
WorkspaceOwnerID: uuid.New(),
WorkspaceID: uuid.New(),
WorkspaceName: "test-workspace",
AgentName: "test-agent",
Type: database.ConnectionTypeWorkspaceApp,
ConnectionID: uuid.NullUUID{},
ConnectionStatus: database.ConnectionStatusConnected,
}
}
@@ -0,0 +1,371 @@
package connectionlog_test
import (
"database/sql"
"net"
"testing"
"time"
"github.com/google/uuid"
"github.com/sqlc-dev/pqtype"
"github.com/stretchr/testify/require"
"cdr.dev/slog/v3"
"cdr.dev/slog/v3/sloggers/slogtest"
"github.com/coder/coder/v2/coderd/database"
"github.com/coder/coder/v2/coderd/database/dbauthz"
"github.com/coder/coder/v2/coderd/database/dbgen"
"github.com/coder/coder/v2/coderd/database/dbtestutil"
"github.com/coder/coder/v2/coderd/database/dbtime"
"github.com/coder/coder/v2/enterprise/coderd/connectionlog"
"github.com/coder/coder/v2/testutil"
"github.com/coder/quartz"
)
func createWorkspace(t *testing.T, db database.Store) database.WorkspaceTable {
t.Helper()
u := dbgen.User(t, db, database.User{})
o := dbgen.Organization(t, db, database.Organization{})
tpl := dbgen.Template(t, db, database.Template{
OrganizationID: o.ID,
CreatedBy: u.ID,
})
return dbgen.Workspace(t, db, database.WorkspaceTable{
ID: uuid.New(),
OwnerID: u.ID,
OrganizationID: o.ID,
AutomaticUpdates: database.AutomaticUpdatesNever,
TemplateID: tpl.ID,
})
}
func testIP() pqtype.Inet {
return pqtype.Inet{
IPNet: net.IPNet{
IP: net.IPv4(127, 0, 0, 1),
Mask: net.IPv4Mask(255, 255, 255, 255),
},
Valid: true,
}
}
func TestDBBackendIntegration(t *testing.T) {
t.Parallel()
t.Run("SingleConnect", func(t *testing.T) {
t.Parallel()
db, _ := dbtestutil.NewDB(t)
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
clock := quartz.NewMock(t)
ws := createWorkspace(t, db)
//nolint:gocritic // Test needs system context for the batcher.
backend := connectionlog.NewDBBatcher(
dbauthz.AsConnectionLogger(ctx), db, log,
connectionlog.WithClock(clock),
connectionlog.WithBatchSize(100),
)
connID := uuid.New()
connectTime := dbtime.Now()
err := backend.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: connectTime,
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: connID, Valid: true},
ConnectionStatus: database.ConnectionStatusConnected,
IP: testIP(),
})
require.NoError(t, err)
err = backend.Close()
require.NoError(t, err)
rows, err := db.GetConnectionLogsOffset(ctx, database.GetConnectionLogsOffsetParams{
LimitOpt: 10,
})
require.NoError(t, err)
require.Len(t, rows, 1)
require.Equal(t, connID, rows[0].ConnectionLog.ConnectionID.UUID)
require.False(t, rows[0].ConnectionLog.DisconnectTime.Valid)
})
t.Run("ConnectThenDisconnectSeparateBatches", func(t *testing.T) {
t.Parallel()
db, _ := dbtestutil.NewDB(t)
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
clock := quartz.NewMock(t)
ws := createWorkspace(t, db)
connID := uuid.New()
connectTime := dbtime.Now()
// First batcher: insert connect, close to flush.
//nolint:gocritic // Test needs system context for the batcher.
b1 := connectionlog.NewDBBatcher(
dbauthz.AsConnectionLogger(ctx), db, log,
connectionlog.WithClock(clock),
connectionlog.WithBatchSize(100),
)
err := b1.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: connectTime,
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: connID, Valid: true},
ConnectionStatus: database.ConnectionStatusConnected,
IP: testIP(),
})
require.NoError(t, err)
require.NoError(t, b1.Close())
// Second batcher: insert disconnect, close to flush.
//nolint:gocritic // Test needs system context for the batcher.
b2 := connectionlog.NewDBBatcher(
dbauthz.AsConnectionLogger(ctx), db, log,
connectionlog.WithClock(clock),
connectionlog.WithBatchSize(100),
)
disconnectTime := connectTime.Add(5 * time.Second)
err = b2.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: disconnectTime,
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: connID, Valid: true},
ConnectionStatus: database.ConnectionStatusDisconnected,
Code: sql.NullInt32{Int32: 0, Valid: true},
DisconnectReason: sql.NullString{String: "client left", Valid: true},
IP: testIP(),
})
require.NoError(t, err)
require.NoError(t, b2.Close())
rows, err := db.GetConnectionLogsOffset(ctx, database.GetConnectionLogsOffsetParams{
LimitOpt: 10,
})
require.NoError(t, err)
require.Len(t, rows, 1, "connect+disconnect should produce one row")
require.True(t, rows[0].ConnectionLog.DisconnectTime.Valid)
require.Equal(t, "client left", rows[0].ConnectionLog.DisconnectReason.String)
})
t.Run("ConnectAndDisconnectSameBatch", func(t *testing.T) {
t.Parallel()
db, _ := dbtestutil.NewDB(t)
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
clock := quartz.NewMock(t)
ws := createWorkspace(t, db)
//nolint:gocritic // Test needs system context for the batcher.
backend := connectionlog.NewDBBatcher(
dbauthz.AsConnectionLogger(ctx), db, log,
connectionlog.WithClock(clock),
connectionlog.WithBatchSize(100),
)
connID := uuid.New()
connectTime := dbtime.Now()
disconnectTime := connectTime.Add(time.Second)
// Both events in the same batch window.
err := backend.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: connectTime,
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: connID, Valid: true},
ConnectionStatus: database.ConnectionStatusConnected,
IP: testIP(),
})
require.NoError(t, err)
err = backend.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: disconnectTime,
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: connID, Valid: true},
ConnectionStatus: database.ConnectionStatusDisconnected,
Code: sql.NullInt32{Int32: 0, Valid: true},
DisconnectReason: sql.NullString{String: "done", Valid: true},
IP: testIP(),
})
require.NoError(t, err)
// Close drains channel and flushes — dedup keeps disconnect.
err = backend.Close()
require.NoError(t, err)
rows, err := db.GetConnectionLogsOffset(ctx, database.GetConnectionLogsOffsetParams{
LimitOpt: 10,
})
require.NoError(t, err)
require.Len(t, rows, 1)
require.True(t, rows[0].ConnectionLog.DisconnectTime.Valid)
require.Equal(t, "done", rows[0].ConnectionLog.DisconnectReason.String)
})
t.Run("MultipleIndependentConnections", func(t *testing.T) {
t.Parallel()
db, _ := dbtestutil.NewDB(t)
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
clock := quartz.NewMock(t)
ws := createWorkspace(t, db)
//nolint:gocritic // Test needs system context for the batcher.
backend := connectionlog.NewDBBatcher(
dbauthz.AsConnectionLogger(ctx), db, log,
connectionlog.WithClock(clock),
connectionlog.WithBatchSize(100),
)
now := dbtime.Now()
for i := 0; i < 5; i++ {
err := backend.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: now,
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: uuid.New(), Valid: true},
ConnectionStatus: database.ConnectionStatusConnected,
IP: testIP(),
})
require.NoError(t, err)
}
err := backend.Close()
require.NoError(t, err)
rows, err := db.GetConnectionLogsOffset(ctx, database.GetConnectionLogsOffsetParams{
LimitOpt: 10,
})
require.NoError(t, err)
require.Len(t, rows, 5)
})
t.Run("NullConnectionIDWebEvents", func(t *testing.T) {
t.Parallel()
db, _ := dbtestutil.NewDB(t)
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
clock := quartz.NewMock(t)
ws := createWorkspace(t, db)
//nolint:gocritic // Test needs system context for the batcher.
backend := connectionlog.NewDBBatcher(
dbauthz.AsConnectionLogger(ctx), db, log,
connectionlog.WithClock(clock),
connectionlog.WithBatchSize(100),
)
now := dbtime.Now()
for i := 0; i < 2; i++ {
err := backend.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: now,
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeWorkspaceApp,
ConnectionID: uuid.NullUUID{},
ConnectionStatus: database.ConnectionStatusConnected,
IP: testIP(),
})
require.NoError(t, err)
}
err := backend.Close()
require.NoError(t, err)
rows, err := db.GetConnectionLogsOffset(ctx, database.GetConnectionLogsOffsetParams{
LimitOpt: 10,
})
require.NoError(t, err)
require.Len(t, rows, 2, "null connection_id events should not be deduplicated")
})
t.Run("CloseFlushesToDB", func(t *testing.T) {
t.Parallel()
db, _ := dbtestutil.NewDB(t)
ctx := testutil.Context(t, testutil.WaitShort)
log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug)
clock := quartz.NewMock(t)
ws := createWorkspace(t, db)
//nolint:gocritic // Test needs system context for the batcher.
backend := connectionlog.NewDBBatcher(
dbauthz.AsConnectionLogger(ctx), db, log,
connectionlog.WithClock(clock),
connectionlog.WithBatchSize(100),
)
err := backend.Upsert(ctx, database.UpsertConnectionLogParams{
ID: uuid.New(),
Time: dbtime.Now(),
OrganizationID: ws.OrganizationID,
WorkspaceOwnerID: ws.OwnerID,
WorkspaceID: ws.ID,
WorkspaceName: ws.Name,
AgentName: "main",
Type: database.ConnectionTypeSsh,
ConnectionID: uuid.NullUUID{UUID: uuid.New(), Valid: true},
ConnectionStatus: database.ConnectionStatusConnected,
IP: testIP(),
})
require.NoError(t, err)
// Close without advancing clock — final flush should write.
err = backend.Close()
require.NoError(t, err)
rows, err := db.GetConnectionLogsOffset(ctx, database.GetConnectionLogsOffsetParams{
LimitOpt: 10,
})
require.NoError(t, err)
require.Len(t, rows, 1)
})
}
+1 -1
View File
@@ -227,7 +227,7 @@ func TestConnectionLogs(t *testing.T) {
Int32: 0,
Valid: false,
},
Ip: pqtype.Inet{IPNet: net.IPNet{
IP: pqtype.Inet{IPNet: net.IPNet{
IP: net.ParseIP("192.168.0.1"),
Mask: net.CIDRMask(8, 32),
}, Valid: true},
+3 -3
View File
@@ -784,7 +784,7 @@ func TestIssueSignedAppToken(t *testing.T) {
require.NoError(t, err)
require.True(t, connectionLogger.Contains(t, database.UpsertConnectionLogParams{
Ip: parsedFakeClientIP,
IP: parsedFakeClientIP,
}))
})
@@ -812,7 +812,7 @@ func TestIssueSignedAppToken(t *testing.T) {
}
require.True(t, connectionLogger.Contains(t, database.UpsertConnectionLogParams{
Ip: parsedFakeClientIP,
IP: parsedFakeClientIP,
}))
})
}
@@ -1020,7 +1020,7 @@ func TestReconnectingPTYSignedToken(t *testing.T) {
// validate it here.
require.True(t, connectionLogger.Contains(t, database.UpsertConnectionLogParams{
Ip: pqtype.Inet{
IP: pqtype.Inet{
Valid: true, IPNet: net.IPNet{
IP: net.ParseIP("127.0.0.1"),
Mask: net.CIDRMask(32, 32),