mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: implement boundary usage tracker and telemetry collection (#21716)
Implements telemetry for boundary usage tracking across all Coder replicas and reports them via telemetry. Changes: - Implement Tracker with Track(), FlushToDB(), and StartFlushLoop() methods - Add telemetry integration via collectBoundaryUsageSummary() - Use telemetry lock to ensure only one replica collects per period The tracker accumulates unique workspaces, unique users, and request counts (allowed/denied) in memory, then flushes to the database periodically. During telemetry collection, stats are aggregated across all replicas and reset for the next period.
This commit is contained in:
@@ -31,6 +31,7 @@ import (
|
||||
"github.com/coder/coder/v2/buildinfo"
|
||||
clitelemetry "github.com/coder/coder/v2/cli/telemetry"
|
||||
"github.com/coder/coder/v2/coderd/database"
|
||||
"github.com/coder/coder/v2/coderd/database/dbauthz"
|
||||
"github.com/coder/coder/v2/coderd/database/dbtime"
|
||||
"github.com/coder/coder/v2/coderd/util/ptr"
|
||||
"github.com/coder/coder/v2/codersdk"
|
||||
@@ -759,6 +760,17 @@ func (r *remoteReporter) createSnapshot() (*Snapshot, error) {
|
||||
snapshot.AIBridgeInterceptionsSummaries = summaries
|
||||
return nil
|
||||
})
|
||||
eg.Go(func() error {
|
||||
summary, err := r.collectBoundaryUsageSummary(ctx)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("collect boundary usage summary: %w", err)
|
||||
}
|
||||
// Only send a summary if there was actual usage.
|
||||
if summary != nil && summary.UniqueUsers > 0 {
|
||||
snapshot.BoundaryUsageSummary = summary
|
||||
}
|
||||
return nil
|
||||
})
|
||||
|
||||
err := eg.Wait()
|
||||
if err != nil {
|
||||
@@ -837,6 +849,51 @@ func (r *remoteReporter) generateAIBridgeInterceptionsSummaries(ctx context.Cont
|
||||
return summaries, eg.Wait()
|
||||
}
|
||||
|
||||
// collectBoundaryUsageSummary collects boundary usage statistics from all
|
||||
// replicas and resets the stats for the next telemetry period. Returns nil if
|
||||
// another replica has already collected for this period.
|
||||
func (r *remoteReporter) collectBoundaryUsageSummary(ctx context.Context) (*BoundaryUsageSummary, error) {
|
||||
// Use twice the snapshot frequency as the staleness limit to ensure we
|
||||
// capture data from replicas that may have slightly different flush times.
|
||||
maxStaleness := r.options.SnapshotFrequency * 2
|
||||
//nolint:gocritic // This is the actual collection of boundary usage tracking.
|
||||
boundaryCtx := dbauthz.AsBoundaryUsageTracker(ctx)
|
||||
|
||||
// Claim the telemetry lock for this period. Use snapshot frequency so each
|
||||
// telemetry snapshot period gets exactly one collection.
|
||||
now := dbtime.Time(r.options.Clock.Now()).UTC()
|
||||
periodEndingAt := now.Truncate(r.options.SnapshotFrequency)
|
||||
err := r.options.Database.InsertTelemetryLock(ctx, database.InsertTelemetryLockParams{
|
||||
EventType: "boundary_usage_summary",
|
||||
PeriodEndingAt: periodEndingAt,
|
||||
})
|
||||
if database.IsUniqueViolation(err, database.UniqueTelemetryLocksPkey) {
|
||||
r.options.Logger.Debug(ctx, "boundary usage telemetry lock already claimed by another replica, skipping", slog.F("period_ending_at", periodEndingAt))
|
||||
return nil, nil //nolint:nilnil // This is simple to handle when dealing with telemetry.
|
||||
}
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("insert boundary usage telemetry lock (period_ending_at=%q): %w", periodEndingAt, err)
|
||||
}
|
||||
|
||||
summary, err := r.options.Database.GetBoundaryUsageSummary(boundaryCtx, maxStaleness.Milliseconds())
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("get boundary usage summary: %w", err)
|
||||
}
|
||||
|
||||
// Reset stats after capturing the summary. This deletes all rows so each
|
||||
// replica will detect a new period on their next flush.
|
||||
if err := r.options.Database.ResetBoundaryUsageStats(boundaryCtx); err != nil {
|
||||
return nil, xerrors.Errorf("reset boundary usage stats: %w", err)
|
||||
}
|
||||
|
||||
return &BoundaryUsageSummary{
|
||||
UniqueWorkspaces: summary.UniqueWorkspaces,
|
||||
UniqueUsers: summary.UniqueUsers,
|
||||
AllowedRequests: summary.AllowedRequests,
|
||||
DeniedRequests: summary.DeniedRequests,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// ConvertAPIKey anonymizes an API key.
|
||||
func ConvertAPIKey(apiKey database.APIKey) APIKey {
|
||||
a := APIKey{
|
||||
@@ -1309,6 +1366,7 @@ type Snapshot struct {
|
||||
UserTailnetConnections []UserTailnetConnection `json:"user_tailnet_connections"`
|
||||
PrebuiltWorkspaces []PrebuiltWorkspace `json:"prebuilt_workspaces"`
|
||||
AIBridgeInterceptionsSummaries []AIBridgeInterceptionsSummary `json:"aibridge_interceptions_summaries"`
|
||||
BoundaryUsageSummary *BoundaryUsageSummary `json:"boundary_usage_summary"`
|
||||
}
|
||||
|
||||
// Deployment contains information about the host running Coder.
|
||||
@@ -1995,6 +2053,15 @@ type AIBridgeInterceptionsSummary struct {
|
||||
InjectedToolCallErrorCount int64 `json:"injected_tool_call_error_count"`
|
||||
}
|
||||
|
||||
// BoundaryUsageSummary contains aggregated boundary usage statistics across all
|
||||
// replicas for the telemetry period.
|
||||
type BoundaryUsageSummary struct {
|
||||
UniqueWorkspaces int64 `json:"unique_workspaces"`
|
||||
UniqueUsers int64 `json:"unique_users"`
|
||||
AllowedRequests int64 `json:"allowed_requests"`
|
||||
DeniedRequests int64 `json:"denied_requests"`
|
||||
}
|
||||
|
||||
func ConvertAIBridgeInterceptionsSummary(endTime time.Time, provider, model, client string, summary database.CalculateAIBridgeInterceptionsTelemetrySummaryRow) AIBridgeInterceptionsSummary {
|
||||
return AIBridgeInterceptionsSummary{
|
||||
ID: uuid.New(),
|
||||
|
||||
@@ -19,7 +19,9 @@ import (
|
||||
"go.uber.org/goleak"
|
||||
|
||||
"github.com/coder/coder/v2/buildinfo"
|
||||
"github.com/coder/coder/v2/coderd/boundaryusage"
|
||||
"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"
|
||||
@@ -841,3 +843,128 @@ func collectSnapshot(
|
||||
|
||||
return testutil.RequireReceive(ctx, t, deployment), testutil.RequireReceive(ctx, t, snapshot)
|
||||
}
|
||||
|
||||
func TestTelemetry_BoundaryUsageSummary(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
t.Run("IncludedInSnapshot", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db, _ := dbtestutil.NewDB(t)
|
||||
ctx := testutil.Context(t, testutil.WaitMedium)
|
||||
|
||||
tracker := boundaryusage.NewTracker()
|
||||
workspace1, workspace2 := uuid.New(), uuid.New()
|
||||
user1, user2 := uuid.New(), uuid.New()
|
||||
replicaID := uuid.New()
|
||||
|
||||
tracker.Track(workspace1, user1, 10, 2)
|
||||
tracker.Track(workspace2, user1, 5, 1)
|
||||
tracker.Track(workspace2, user2, 3, 0)
|
||||
|
||||
// Flush the tracker to the database.
|
||||
err := tracker.FlushToDB(ctx, db, replicaID)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Collect a snapshot and verify boundary usage is included.
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(dbtime.Now())
|
||||
|
||||
_, snapshot := collectSnapshot(ctx, t, db, func(opts telemetry.Options) telemetry.Options {
|
||||
opts.Clock = clock
|
||||
return opts
|
||||
})
|
||||
|
||||
require.NotNil(t, snapshot.BoundaryUsageSummary)
|
||||
require.Equal(t, int64(2), snapshot.BoundaryUsageSummary.UniqueWorkspaces)
|
||||
require.Equal(t, int64(2), snapshot.BoundaryUsageSummary.UniqueUsers)
|
||||
require.Equal(t, int64(10+5+3), snapshot.BoundaryUsageSummary.AllowedRequests)
|
||||
require.Equal(t, int64(2+1+0), snapshot.BoundaryUsageSummary.DeniedRequests)
|
||||
})
|
||||
|
||||
t.Run("ResetAfterCollection", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db, _ := dbtestutil.NewDB(t)
|
||||
ctx := testutil.Context(t, testutil.WaitMedium)
|
||||
|
||||
tracker := boundaryusage.NewTracker()
|
||||
replicaID := uuid.New()
|
||||
|
||||
tracker.Track(uuid.New(), uuid.New(), 5, 1)
|
||||
err := tracker.FlushToDB(ctx, db, replicaID)
|
||||
require.NoError(t, err)
|
||||
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(dbtime.Now())
|
||||
|
||||
// First snapshot should have the data.
|
||||
_, snapshot1 := collectSnapshot(ctx, t, db, func(opts telemetry.Options) telemetry.Options {
|
||||
opts.Clock = clock
|
||||
return opts
|
||||
})
|
||||
require.NotNil(t, snapshot1.BoundaryUsageSummary)
|
||||
require.Equal(t, int64(5), snapshot1.BoundaryUsageSummary.AllowedRequests)
|
||||
|
||||
// Advance clock to next snapshot period to avoid lock conflict.
|
||||
clock.Advance(30 * time.Minute)
|
||||
|
||||
// Second snapshot should have no data (stats were reset).
|
||||
_, snapshot2 := collectSnapshot(ctx, t, db, func(opts telemetry.Options) telemetry.Options {
|
||||
opts.Clock = clock
|
||||
return opts
|
||||
})
|
||||
// Summary should be nil or have zero values since stats were reset.
|
||||
if snapshot2.BoundaryUsageSummary != nil {
|
||||
require.Equal(t, int64(0), snapshot2.BoundaryUsageSummary.AllowedRequests)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("OnlyOneReplicaCollects", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db, _ := dbtestutil.NewDB(t)
|
||||
ctx := testutil.Context(t, testutil.WaitMedium)
|
||||
|
||||
// Set up boundary usage stats from two replicas.
|
||||
tracker1 := boundaryusage.NewTracker()
|
||||
tracker2 := boundaryusage.NewTracker()
|
||||
replica1ID := uuid.New()
|
||||
replica2ID := uuid.New()
|
||||
|
||||
tracker1.Track(uuid.New(), uuid.New(), 10, 1)
|
||||
tracker2.Track(uuid.New(), uuid.New(), 20, 2)
|
||||
|
||||
err := tracker1.FlushToDB(ctx, db, replica1ID)
|
||||
require.NoError(t, err)
|
||||
err = tracker2.FlushToDB(ctx, db, replica2ID)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Verify both replicas' data is in the database.
|
||||
boundaryCtx := dbauthz.AsBoundaryUsageTracker(ctx)
|
||||
const maxStalenessMs = 60000
|
||||
summary, err := db.GetBoundaryUsageSummary(boundaryCtx, maxStalenessMs)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(10+20), summary.AllowedRequests)
|
||||
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(dbtime.Now())
|
||||
|
||||
// First snapshot collects and resets.
|
||||
_, snapshot1 := collectSnapshot(ctx, t, db, func(opts telemetry.Options) telemetry.Options {
|
||||
opts.Clock = clock
|
||||
return opts
|
||||
})
|
||||
require.NotNil(t, snapshot1.BoundaryUsageSummary)
|
||||
require.Equal(t, int64(10+20), snapshot1.BoundaryUsageSummary.AllowedRequests)
|
||||
|
||||
// Second snapshot in same period should skip (lock already claimed).
|
||||
_, snapshot2 := collectSnapshot(ctx, t, db, func(opts telemetry.Options) telemetry.Options {
|
||||
opts.Clock = clock
|
||||
return opts
|
||||
})
|
||||
// The second snapshot should have nil because another "replica" already
|
||||
// claimed the lock for this period.
|
||||
require.Nil(t, snapshot2.BoundaryUsageSummary)
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user