From 102d8601573c78d603a31309d6ec0d74acc35900 Mon Sep 17 00:00:00 2001 From: Ehab Younes Date: Thu, 13 Aug 2026 18:09:24 +0300 Subject: [PATCH] perf: store workspace agent session counts as JSONB --- coderd/database/dbgen/dbgen.go | 81 ++-- coderd/database/dbpurge/dbpurge_test.go | 9 +- coderd/database/dbrollup/dbrollup_test.go | 40 +- coderd/database/dump.sql | 11 +- ...69_workspace_agent_session_counts.down.sql | 22 + ...0569_workspace_agent_session_counts.up.sql | 70 ++++ coderd/database/migrations/migrate_test.go | 153 +++++++ ...0569_workspace_agent_session_counts.up.sql | 8 + coderd/database/models.go | 34 +- coderd/database/querier.go | 1 - coderd/database/querier_test.go | 332 +++++++++------- coderd/database/queries.sql.go | 376 ++++++++---------- coderd/database/queries/insights.sql | 86 ++-- .../database/queries/workspaceagentstats.sql | 244 +++++------- coderd/metricscache/metricscache_test.go | 30 +- coderd/workspacestats/batcher.go | 71 ++-- .../workspacestats/batcher_internal_test.go | 52 ++- 17 files changed, 910 insertions(+), 710 deletions(-) create mode 100644 coderd/database/migrations/000569_workspace_agent_session_counts.down.sql create mode 100644 coderd/database/migrations/000569_workspace_agent_session_counts.up.sql create mode 100644 coderd/database/migrations/testdata/fixtures/000569_workspace_agent_session_counts.up.sql diff --git a/coderd/database/dbgen/dbgen.go b/coderd/database/dbgen/dbgen.go index 10cfe4dddff08..c4a8193a44e0d 100644 --- a/coderd/database/dbgen/dbgen.go +++ b/coderd/database/dbgen/dbgen.go @@ -1673,54 +1673,57 @@ func TemplateVersionTerraformValues(t testing.TB, db database.Store, orig databa return v } -func WorkspaceAgentStat(t testing.TB, db database.Store, orig database.WorkspaceAgentStat) database.WorkspaceAgentStat { +// WorkspaceAgentStat inserts a workspace agent stat row. The optional map seeds +// its session_counts column. +func WorkspaceAgentStat(t testing.TB, db database.Store, orig database.WorkspaceAgentStat, sessionCounts ...map[string]int64) database.WorkspaceAgentStat { if orig.ConnectionsByProto == nil { orig.ConnectionsByProto = json.RawMessage([]byte("{}")) } jsonProto := []byte(fmt.Sprintf("[%s]", orig.ConnectionsByProto)) + // The insert rejects null session count elements. + counts := map[string]int64{} + if len(sessionCounts) > 0 && sessionCounts[0] != nil { + counts = sessionCounts[0] + } + jsonSessionCounts, err := json.Marshal([]map[string]int64{counts}) + require.NoError(t, err, "marshal session counts") + params := database.InsertWorkspaceAgentStatsParams{ - ID: []uuid.UUID{takeFirst(orig.ID, uuid.New())}, - CreatedAt: []time.Time{takeFirst(orig.CreatedAt, dbtime.Now())}, - UserID: []uuid.UUID{takeFirst(orig.UserID, uuid.New())}, - TemplateID: []uuid.UUID{takeFirst(orig.TemplateID, uuid.New())}, - WorkspaceID: []uuid.UUID{takeFirst(orig.WorkspaceID, uuid.New())}, - AgentID: []uuid.UUID{takeFirst(orig.AgentID, uuid.New())}, - ConnectionsByProto: jsonProto, - ConnectionCount: []int64{takeFirst(orig.ConnectionCount, 0)}, - RxPackets: []int64{takeFirst(orig.RxPackets, 0)}, - RxBytes: []int64{takeFirst(orig.RxBytes, 0)}, - TxPackets: []int64{takeFirst(orig.TxPackets, 0)}, - TxBytes: []int64{takeFirst(orig.TxBytes, 0)}, - SessionCountVSCode: []int64{takeFirst(orig.SessionCountVSCode, 0)}, - SessionCountJetBrains: []int64{takeFirst(orig.SessionCountJetBrains, 0)}, - SessionCountReconnectingPTY: []int64{takeFirst(orig.SessionCountReconnectingPTY, 0)}, - SessionCountSSH: []int64{takeFirst(orig.SessionCountSSH, 0)}, - ConnectionMedianLatencyMS: []float64{takeFirst(orig.ConnectionMedianLatencyMS, 0)}, - Usage: []bool{takeFirst(orig.Usage, false)}, - } - err := db.InsertWorkspaceAgentStats(genCtx, params) + ID: []uuid.UUID{takeFirst(orig.ID, uuid.New())}, + CreatedAt: []time.Time{takeFirst(orig.CreatedAt, dbtime.Now())}, + UserID: []uuid.UUID{takeFirst(orig.UserID, uuid.New())}, + TemplateID: []uuid.UUID{takeFirst(orig.TemplateID, uuid.New())}, + WorkspaceID: []uuid.UUID{takeFirst(orig.WorkspaceID, uuid.New())}, + AgentID: []uuid.UUID{takeFirst(orig.AgentID, uuid.New())}, + ConnectionsByProto: jsonProto, + ConnectionCount: []int64{takeFirst(orig.ConnectionCount, 0)}, + RxPackets: []int64{takeFirst(orig.RxPackets, 0)}, + RxBytes: []int64{takeFirst(orig.RxBytes, 0)}, + TxPackets: []int64{takeFirst(orig.TxPackets, 0)}, + TxBytes: []int64{takeFirst(orig.TxBytes, 0)}, + SessionCounts: jsonSessionCounts, + ConnectionMedianLatencyMS: []float64{takeFirst(orig.ConnectionMedianLatencyMS, 0)}, + Usage: []bool{takeFirst(orig.Usage, false)}, + } + err = db.InsertWorkspaceAgentStats(genCtx, params) require.NoError(t, err, "insert workspace agent stat") return database.WorkspaceAgentStat{ - ID: params.ID[0], - CreatedAt: params.CreatedAt[0], - UserID: params.UserID[0], - AgentID: params.AgentID[0], - WorkspaceID: params.WorkspaceID[0], - TemplateID: params.TemplateID[0], - ConnectionsByProto: orig.ConnectionsByProto, - ConnectionCount: params.ConnectionCount[0], - RxPackets: params.RxPackets[0], - RxBytes: params.RxBytes[0], - TxPackets: params.TxPackets[0], - TxBytes: params.TxBytes[0], - ConnectionMedianLatencyMS: params.ConnectionMedianLatencyMS[0], - SessionCountVSCode: params.SessionCountVSCode[0], - SessionCountJetBrains: params.SessionCountJetBrains[0], - SessionCountReconnectingPTY: params.SessionCountReconnectingPTY[0], - SessionCountSSH: params.SessionCountSSH[0], - Usage: params.Usage[0], + ID: params.ID[0], + CreatedAt: params.CreatedAt[0], + UserID: params.UserID[0], + AgentID: params.AgentID[0], + WorkspaceID: params.WorkspaceID[0], + TemplateID: params.TemplateID[0], + ConnectionsByProto: orig.ConnectionsByProto, + ConnectionCount: params.ConnectionCount[0], + RxPackets: params.RxPackets[0], + RxBytes: params.RxBytes[0], + TxPackets: params.TxPackets[0], + TxBytes: params.TxBytes[0], + ConnectionMedianLatencyMS: params.ConnectionMedianLatencyMS[0], + Usage: params.Usage[0], } } diff --git a/coderd/database/dbpurge/dbpurge_test.go b/coderd/database/dbpurge/dbpurge_test.go index 242f9c33324ee..c41023cb62c56 100644 --- a/coderd/database/dbpurge/dbpurge_test.go +++ b/coderd/database/dbpurge/dbpurge_test.go @@ -387,8 +387,7 @@ func TestDeleteOldWorkspaceAgentStats(t *testing.T) { ConnectionCount: 1, ConnectionMedianLatencyMS: 1, RxBytes: 1111, - SessionCountSSH: 1, - }) + }, map[string]int64{"ssh": 1}) // Stat inserted 180 days - 2 hour ago, should not be deleted before rollup. second := dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ @@ -396,8 +395,7 @@ func TestDeleteOldWorkspaceAgentStats(t *testing.T) { ConnectionCount: 1, ConnectionMedianLatencyMS: 1, RxBytes: 2222, - SessionCountSSH: 1, - }) + }, map[string]int64{"ssh": 1}) // Stat inserted 179 days - 4 hour ago, should not be deleted at all. third := dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ @@ -405,8 +403,7 @@ func TestDeleteOldWorkspaceAgentStats(t *testing.T) { ConnectionCount: 1, ConnectionMedianLatencyMS: 1, RxBytes: 3333, - SessionCountSSH: 1, - }) + }, map[string]int64{"ssh": 1}) // when closer := dbpurge.New(ctx, logger, db, &codersdk.DeploymentValues{}, prometheus.NewRegistry(), dbpurge.WithClock(clk)) diff --git a/coderd/database/dbrollup/dbrollup_test.go b/coderd/database/dbrollup/dbrollup_test.go index ebcb6852a3fcb..4be76c89279c5 100644 --- a/coderd/database/dbrollup/dbrollup_test.go +++ b/coderd/database/dbrollup/dbrollup_test.go @@ -75,8 +75,7 @@ func TestRollup_TwoInstancesUseLocking(t *testing.T) { CreatedAt: refTime.Add(-time.Minute), ConnectionMedianLatencyMS: 1, ConnectionCount: 1, - SessionCountSSH: 1, - }) + }, map[string]int64{"ssh": 1}) closeRolluper := func(rolluper *dbrollup.Rolluper, resume chan struct{}) { close(resume) @@ -163,8 +162,7 @@ func TestRollupTemplateUsageStats(t *testing.T) { CreatedAt: anHourAndSixMonthsAgo.AddDate(0, 0, -1), ConnectionMedianLatencyMS: 1, ConnectionCount: 1, - SessionCountSSH: 1, - }) + }, map[string]int64{"ssh": 1}) _ = dbgen.WorkspaceAppStat(t, db, database.WorkspaceAppStat{ UserID: user.ID, WorkspaceID: ws.ID, @@ -176,25 +174,23 @@ func TestRollupTemplateUsageStats(t *testing.T) { // Stats inserted 6 months - 1 day ago, should be rolled up. wags1 := dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - TemplateID: tpl.ID, - WorkspaceID: ws.ID, - AgentID: agent.ID, - UserID: user.ID, - CreatedAt: anHourAndSixMonthsAgo.AddDate(0, 0, 1), - ConnectionMedianLatencyMS: 1, - ConnectionCount: 1, - SessionCountReconnectingPTY: 1, - }) + TemplateID: tpl.ID, + WorkspaceID: ws.ID, + AgentID: agent.ID, + UserID: user.ID, + CreatedAt: anHourAndSixMonthsAgo.AddDate(0, 0, 1), + ConnectionMedianLatencyMS: 1, + ConnectionCount: 1, + }, map[string]int64{"reconnecting_pty": 1}) wags2 := dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - TemplateID: tpl.ID, - WorkspaceID: ws.ID, - AgentID: agent.ID, - UserID: user.ID, - CreatedAt: wags1.CreatedAt.Add(time.Minute), - ConnectionMedianLatencyMS: 1, - ConnectionCount: 1, - SessionCountReconnectingPTY: 1, - }) + TemplateID: tpl.ID, + WorkspaceID: ws.ID, + AgentID: agent.ID, + UserID: user.ID, + CreatedAt: wags1.CreatedAt.Add(time.Minute), + ConnectionMedianLatencyMS: 1, + ConnectionCount: 1, + }, map[string]int64{"reconnecting_pty": 1}) // wags2 and waps1 overlap, so total usage is 4 - 1. waps1 := dbgen.WorkspaceAppStat(t, db, database.WorkspaceAppStat{ UserID: user.ID, diff --git a/coderd/database/dump.sql b/coderd/database/dump.sql index 79e987cc084e9..65c610385b734 100644 --- a/coderd/database/dump.sql +++ b/coderd/database/dump.sql @@ -3851,13 +3851,12 @@ CREATE TABLE workspace_agent_stats ( tx_packets bigint DEFAULT 0 NOT NULL, tx_bytes bigint DEFAULT 0 NOT NULL, connection_median_latency_ms double precision DEFAULT '-1'::integer NOT NULL, - session_count_vscode bigint DEFAULT 0 NOT NULL, - session_count_jetbrains bigint DEFAULT 0 NOT NULL, - session_count_reconnecting_pty bigint DEFAULT 0 NOT NULL, - session_count_ssh bigint DEFAULT 0 NOT NULL, - usage boolean DEFAULT false NOT NULL + usage boolean DEFAULT false NOT NULL, + session_counts jsonb DEFAULT '{}'::jsonb NOT NULL ); +COMMENT ON COLUMN workspace_agent_stats.session_counts IS 'Positive session counts keyed by the canonical app name reported by the agent.'; + CREATE TABLE workspace_agent_volume_resource_monitors ( agent_id uuid NOT NULL, enabled boolean NOT NULL, @@ -4973,7 +4972,7 @@ COMMENT ON INDEX workspace_agent_scripts_workspace_agent_id_idx IS 'Foreign key CREATE INDEX workspace_agent_startup_logs_id_agent_id_idx ON workspace_agent_logs USING btree (agent_id, id); -CREATE INDEX workspace_agent_stats_template_id_created_at_user_id_idx ON workspace_agent_stats USING btree (template_id, created_at, user_id) INCLUDE (session_count_vscode, session_count_jetbrains, session_count_reconnecting_pty, session_count_ssh, connection_median_latency_ms) WHERE (connection_count > 0); +CREATE INDEX workspace_agent_stats_template_id_created_at_user_id_idx ON workspace_agent_stats USING btree (template_id, created_at, user_id) INCLUDE (connection_median_latency_ms) WHERE (connection_count > 0); COMMENT ON INDEX workspace_agent_stats_template_id_created_at_user_id_idx IS 'Support index for template insights endpoint to build interval reports faster.'; diff --git a/coderd/database/migrations/000569_workspace_agent_session_counts.down.sql b/coderd/database/migrations/000569_workspace_agent_session_counts.down.sql new file mode 100644 index 0000000000000..ed2d8ed4500dc --- /dev/null +++ b/coderd/database/migrations/000569_workspace_agent_session_counts.down.sql @@ -0,0 +1,22 @@ +ALTER TABLE workspace_agent_stats + ADD COLUMN session_count_vscode bigint DEFAULT 0 NOT NULL, + ADD COLUMN session_count_jetbrains bigint DEFAULT 0 NOT NULL, + ADD COLUMN session_count_reconnecting_pty bigint DEFAULT 0 NOT NULL, + ADD COLUMN session_count_ssh bigint DEFAULT 0 NOT NULL; + +-- Restore the four known session counts. Other keys are discarded. +UPDATE workspace_agent_stats +SET + session_count_vscode = COALESCE((session_counts ->> 'vscode')::bigint, 0), + session_count_jetbrains = COALESCE((session_counts ->> 'jetbrains')::bigint, 0), + session_count_reconnecting_pty = COALESCE((session_counts ->> 'reconnecting_pty')::bigint, 0), + session_count_ssh = COALESCE((session_counts ->> 'ssh')::bigint, 0); + +DROP INDEX workspace_agent_stats_template_id_created_at_user_id_idx; + +ALTER TABLE workspace_agent_stats + DROP COLUMN session_counts; + +CREATE INDEX workspace_agent_stats_template_id_created_at_user_id_idx ON workspace_agent_stats USING btree (template_id, created_at, user_id) INCLUDE (session_count_vscode, session_count_jetbrains, session_count_reconnecting_pty, session_count_ssh, connection_median_latency_ms) WHERE (connection_count > 0); + +COMMENT ON INDEX workspace_agent_stats_template_id_created_at_user_id_idx IS 'Support index for template insights endpoint to build interval reports faster.'; diff --git a/coderd/database/migrations/000569_workspace_agent_session_counts.up.sql b/coderd/database/migrations/000569_workspace_agent_session_counts.up.sql new file mode 100644 index 0000000000000..80e460427b0ed --- /dev/null +++ b/coderd/database/migrations/000569_workspace_agent_session_counts.up.sql @@ -0,0 +1,70 @@ +LOCK TABLE workspace_agent_stats IN ACCESS EXCLUSIVE MODE; + +DO $$ +DECLARE + latest_rollup_start timestamptz; + migration_cutoff timestamptz; +BEGIN + SELECT MAX(start_time) INTO latest_rollup_start FROM template_usage_stats; + migration_cutoff := COALESCE(latest_rollup_start - interval '1 day', statement_timestamp() - interval '180 days'); + + IF (latest_rollup_start IS NULL OR latest_rollup_start < statement_timestamp() - interval '24 hours') + AND EXISTS ( + SELECT 1 + FROM workspace_agent_stats + WHERE created_at >= migration_cutoff + AND ( + session_count_vscode > 0 + OR session_count_jetbrains > 0 + OR session_count_reconnecting_pty > 0 + OR session_count_ssh > 0 + ) + ) + THEN + RAISE EXCEPTION 'migration 000569 requires template usage stats rolled up within the last 24 hours; run the previous Coder version until template usage stats roll up, then retry the upgrade' + USING DETAIL = format( + 'Latest template_usage_stats.start_time: %s.', + COALESCE(latest_rollup_start::text, 'none') + ), + HINT = 'Check coderd logs for "failed to rollup data" if the timestamp does not advance.'; + END IF; +END +$$; + +ALTER TABLE workspace_agent_stats + ADD COLUMN session_counts jsonb DEFAULT '{}'::jsonb NOT NULL; + +COMMENT ON COLUMN workspace_agent_stats.session_counts IS 'Positive session counts keyed by the canonical app name reported by the agent.'; + +-- Convert the raw window still used by rollups and operational statistics. +-- Older usage already lives in template usage rollups. +UPDATE workspace_agent_stats +SET session_counts = jsonb_strip_nulls(jsonb_build_object( + 'vscode', CASE WHEN session_count_vscode > 0 THEN session_count_vscode END, + 'jetbrains', CASE WHEN session_count_jetbrains > 0 THEN session_count_jetbrains END, + 'reconnecting_pty', CASE WHEN session_count_reconnecting_pty > 0 THEN session_count_reconnecting_pty END, + 'ssh', CASE WHEN session_count_ssh > 0 THEN session_count_ssh END +)) +WHERE created_at >= ( + SELECT COALESCE(MAX(start_time) - interval '1 day', statement_timestamp() - interval '180 days') + FROM template_usage_stats +) + AND ( + session_count_vscode > 0 + OR session_count_jetbrains > 0 + OR session_count_reconnecting_pty > 0 + OR session_count_ssh > 0 + ); + +-- Recreate the index without the dropped columns in its INCLUDE list. +DROP INDEX workspace_agent_stats_template_id_created_at_user_id_idx; + +ALTER TABLE workspace_agent_stats + DROP COLUMN session_count_vscode, + DROP COLUMN session_count_jetbrains, + DROP COLUMN session_count_reconnecting_pty, + DROP COLUMN session_count_ssh; + +CREATE INDEX workspace_agent_stats_template_id_created_at_user_id_idx ON workspace_agent_stats USING btree (template_id, created_at, user_id) INCLUDE (connection_median_latency_ms) WHERE (connection_count > 0); + +COMMENT ON INDEX workspace_agent_stats_template_id_created_at_user_id_idx IS 'Support index for template insights endpoint to build interval reports faster.'; diff --git a/coderd/database/migrations/migrate_test.go b/coderd/database/migrations/migrate_test.go index df0c7d14ea9bb..6efbf51a69f88 100644 --- a/coderd/database/migrations/migrate_test.go +++ b/coderd/database/migrations/migrate_test.go @@ -3,6 +3,7 @@ package migrations_test import ( "context" "database/sql" + "errors" "fmt" "os" "path/filepath" @@ -2878,6 +2879,158 @@ func TestMigration000565OAuth2ClientTypeConstraint(t *testing.T) { } } +//nolint:tparallel,paralleltest // Subtests share one database with transaction-local fixtures. +func TestMigration000569WorkspaceAgentSessionCountsGuard(t *testing.T) { + t.Parallel() + + const ( + priorMigrationVersion = 568 + noStatsAgeHours = -1 + ) + + sqlDB := testSQLDB(t) + next, err := migrations.Stepper(sqlDB) + require.NoError(t, err) + for { + version, more, err := next() + require.NoError(t, err) + if !more { + t.Fatalf("migration %d not found", priorMigrationVersion) + } + if version == priorMigrationVersion { + break + } + } + + ctx := testutil.Context(t, testutil.WaitSuperLong) + migrationSQL, err := os.ReadFile("000569_workspace_agent_session_counts.up.sql") + require.NoError(t, err) + + tests := []struct { + name string + statsAgeHours int + sessionCount int + watermarkAgeHours []int + wantError bool + wantCounts string + }{ + {name: "empty database", statsAgeHours: noStatsAgeHours}, + {name: "missing watermark with retained activity", sessionCount: 2, wantError: true}, + {name: "missing watermark with retained idle stats", wantCounts: `{}`}, + {name: "missing watermark with expired activity", statsAgeHours: 181 * 24, sessionCount: 2, wantCounts: `{}`}, + {name: "stale watermark with retained activity", sessionCount: 2, watermarkAgeHours: []int{25}, wantError: true}, + {name: "stale watermark with rolled up activity", statsAgeHours: 50, sessionCount: 2, watermarkAgeHours: []int{25}, wantCounts: `{}`}, + {name: "fresh watermark", sessionCount: 2, watermarkAgeHours: []int{23}, wantCounts: `{"vscode": 2}`}, + {name: "latest watermark is fresh", sessionCount: 2, watermarkAgeHours: []int{25, 23}, wantCounts: `{"vscode": 2}`}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + tx, err := sqlDB.BeginTx(ctx, nil) + require.NoError(t, err) + t.Cleanup(func() { _ = tx.Rollback() }) + + if tt.statsAgeHours != noStatsAgeHours { + _, err = tx.ExecContext(ctx, ` + INSERT INTO workspace_agent_stats ( + id, created_at, user_id, agent_id, workspace_id, template_id, + connection_count, session_count_vscode + ) VALUES ( + gen_random_uuid(), statement_timestamp() - $1::bigint * interval '1 hour', + gen_random_uuid(), gen_random_uuid(), gen_random_uuid(), gen_random_uuid(), 1, $2 + ) + `, tt.statsAgeHours, tt.sessionCount) + require.NoError(t, err) + } + + for _, ageHours := range tt.watermarkAgeHours { + _, err = tx.ExecContext(ctx, ` + INSERT INTO template_usage_stats ( + start_time, end_time, template_id, user_id, median_latency_ms, + usage_mins, ssh_mins, sftp_mins, reconnecting_pty_mins, + vscode_mins, jetbrains_mins, app_usage_mins + ) VALUES ( + statement_timestamp() - $1::bigint * interval '1 hour', + statement_timestamp() - $1::bigint * interval '1 hour' + interval '30 minutes', + gen_random_uuid(), gen_random_uuid(), NULL, 0, 0, 0, 0, 0, 0, NULL + ) + `, ageHours) + require.NoError(t, err) + } + + _, err = tx.ExecContext(ctx, string(migrationSQL)) + if tt.wantError { + require.ErrorContains(t, err, "requires template usage stats rolled up within the last 24 hours") + var pqErr *pq.Error + require.True(t, errors.As(err, &pqErr)) + require.Contains(t, pqErr.Hint, "failed to rollup data") + require.NoError(t, tx.Rollback()) + + var columnExists bool + err = sqlDB.QueryRowContext(ctx, ` + SELECT EXISTS ( + SELECT 1 FROM information_schema.columns + WHERE table_schema = 'public' + AND table_name = 'workspace_agent_stats' + AND column_name = 'session_counts' + ) + `).Scan(&columnExists) + require.NoError(t, err) + require.False(t, columnExists) + return + } + require.NoError(t, err) + + if tt.wantCounts != "" { + var sessionCounts []byte + err = tx.QueryRowContext(ctx, `SELECT session_counts FROM workspace_agent_stats`).Scan(&sessionCounts) + require.NoError(t, err) + require.JSONEq(t, tt.wantCounts, string(sessionCounts)) + } + }) + } + + // The down migration has columns for the four known apps only, so a count + // under any other name is lost. + t.Run("down restores known apps only", func(t *testing.T) { + tx, err := sqlDB.BeginTx(ctx, nil) + require.NoError(t, err) + t.Cleanup(func() { _ = tx.Rollback() }) + + _, err = tx.ExecContext(ctx, string(migrationSQL)) + require.NoError(t, err) + + _, err = tx.ExecContext(ctx, ` + INSERT INTO workspace_agent_stats ( + id, created_at, user_id, agent_id, workspace_id, template_id, + connection_count, session_counts + ) VALUES ( + gen_random_uuid(), statement_timestamp(), gen_random_uuid(), + gen_random_uuid(), gen_random_uuid(), gen_random_uuid(), 1, + '{"vscode": 2, "some_future_ide": 3}'::jsonb + ) + `) + require.NoError(t, err) + + downSQL, err := os.ReadFile("000569_workspace_agent_session_counts.down.sql") + require.NoError(t, err) + _, err = tx.ExecContext(ctx, string(downSQL)) + require.NoError(t, err) + + var vscode, jetbrains, reconnectingPTY, ssh int64 + err = tx.QueryRowContext(ctx, ` + SELECT session_count_vscode, session_count_jetbrains, + session_count_reconnecting_pty, session_count_ssh + FROM workspace_agent_stats + `).Scan(&vscode, &jetbrains, &reconnectingPTY, &ssh) + require.NoError(t, err) + require.EqualValues(t, 2, vscode) + require.EqualValues(t, 0, jetbrains) + require.EqualValues(t, 0, reconnectingPTY) + require.EqualValues(t, 0, ssh) + }) +} + // TestMigration000566OAuth2AuthMethodBackfill covers the repair the backfill // exists for, which the testdata/fixtures run does not reach: its only // oauth2_provider_apps row seeds no token_endpoint_auth_method at all, so the diff --git a/coderd/database/migrations/testdata/fixtures/000569_workspace_agent_session_counts.up.sql b/coderd/database/migrations/testdata/fixtures/000569_workspace_agent_session_counts.up.sql new file mode 100644 index 0000000000000..d2eba9e9bd83b --- /dev/null +++ b/coderd/database/migrations/testdata/fixtures/000569_workspace_agent_session_counts.up.sql @@ -0,0 +1,8 @@ +INSERT INTO workspace_agent_stats ( + id, created_at, user_id, agent_id, workspace_id, template_id, + connection_count, connection_median_latency_ms, session_counts +) VALUES ( + '4a382ba5-6e57-4a58-991e-d4ac4f6c1012', NOW(), + gen_random_uuid(), gen_random_uuid(), gen_random_uuid(), gen_random_uuid(), + 1, 1, '{"vscode": 2, "some_future_ide": 1}'::jsonb +); diff --git a/coderd/database/models.go b/coderd/database/models.go index a68b7e54bc924..e5ccf6dff1979 100644 --- a/coderd/database/models.go +++ b/coderd/database/models.go @@ -6507,24 +6507,22 @@ type WorkspaceAgentScriptTiming struct { } type WorkspaceAgentStat struct { - ID uuid.UUID `db:"id" json:"id"` - CreatedAt time.Time `db:"created_at" json:"created_at"` - UserID uuid.UUID `db:"user_id" json:"user_id"` - AgentID uuid.UUID `db:"agent_id" json:"agent_id"` - WorkspaceID uuid.UUID `db:"workspace_id" json:"workspace_id"` - TemplateID uuid.UUID `db:"template_id" json:"template_id"` - ConnectionsByProto json.RawMessage `db:"connections_by_proto" json:"connections_by_proto"` - ConnectionCount int64 `db:"connection_count" json:"connection_count"` - RxPackets int64 `db:"rx_packets" json:"rx_packets"` - RxBytes int64 `db:"rx_bytes" json:"rx_bytes"` - TxPackets int64 `db:"tx_packets" json:"tx_packets"` - TxBytes int64 `db:"tx_bytes" json:"tx_bytes"` - ConnectionMedianLatencyMS float64 `db:"connection_median_latency_ms" json:"connection_median_latency_ms"` - SessionCountVSCode int64 `db:"session_count_vscode" json:"session_count_vscode"` - SessionCountJetBrains int64 `db:"session_count_jetbrains" json:"session_count_jetbrains"` - SessionCountReconnectingPTY int64 `db:"session_count_reconnecting_pty" json:"session_count_reconnecting_pty"` - SessionCountSSH int64 `db:"session_count_ssh" json:"session_count_ssh"` - Usage bool `db:"usage" json:"usage"` + ID uuid.UUID `db:"id" json:"id"` + CreatedAt time.Time `db:"created_at" json:"created_at"` + UserID uuid.UUID `db:"user_id" json:"user_id"` + AgentID uuid.UUID `db:"agent_id" json:"agent_id"` + WorkspaceID uuid.UUID `db:"workspace_id" json:"workspace_id"` + TemplateID uuid.UUID `db:"template_id" json:"template_id"` + ConnectionsByProto json.RawMessage `db:"connections_by_proto" json:"connections_by_proto"` + ConnectionCount int64 `db:"connection_count" json:"connection_count"` + RxPackets int64 `db:"rx_packets" json:"rx_packets"` + RxBytes int64 `db:"rx_bytes" json:"rx_bytes"` + TxPackets int64 `db:"tx_packets" json:"tx_packets"` + TxBytes int64 `db:"tx_bytes" json:"tx_bytes"` + ConnectionMedianLatencyMS float64 `db:"connection_median_latency_ms" json:"connection_median_latency_ms"` + Usage bool `db:"usage" json:"usage"` + // Positive session counts keyed by the canonical app name reported by the agent. + SessionCounts json.RawMessage `db:"session_counts" json:"session_counts"` } type WorkspaceAgentVolumeResourceMonitor struct { diff --git a/coderd/database/querier.go b/coderd/database/querier.go index bbc90e02859ef..9069dee28600e 100644 --- a/coderd/database/querier.go +++ b/coderd/database/querier.go @@ -975,7 +975,6 @@ type sqlcQuerier interface { GetWorkspaceAgentScriptsByAgentIDs(ctx context.Context, ids []uuid.UUID) ([]GetWorkspaceAgentScriptsByAgentIDsRow, error) GetWorkspaceAgentStats(ctx context.Context, createdAt time.Time) ([]GetWorkspaceAgentStatsRow, error) GetWorkspaceAgentStatsAndLabels(ctx context.Context, createdAt time.Time) ([]GetWorkspaceAgentStatsAndLabelsRow, error) - // `minute_buckets` could return 0 rows if there are no usage stats since `created_at`. GetWorkspaceAgentUsageStats(ctx context.Context, createdAt time.Time) ([]GetWorkspaceAgentUsageStatsRow, error) GetWorkspaceAgentUsageStatsAndLabels(ctx context.Context, createdAt time.Time) ([]GetWorkspaceAgentUsageStatsAndLabelsRow, error) GetWorkspaceAgentsByInstanceID(ctx context.Context, authInstanceID string) ([]WorkspaceAgent, error) diff --git a/coderd/database/querier_test.go b/coderd/database/querier_test.go index 084e61602c843..ff6b8464578a6 100644 --- a/coderd/database/querier_test.go +++ b/coderd/database/querier_test.go @@ -58,14 +58,12 @@ func TestGetDeploymentWorkspaceAgentStats(t *testing.T) { TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 1, - SessionCountVSCode: 1, - }) + }, map[string]int64{"vscode": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 2, - SessionCountVSCode: 1, - }) + }, map[string]int64{"vscode": 1}) stats, err := db.GetDeploymentWorkspaceAgentStats(ctx, dbtime.Now().Add(-time.Hour)) require.NoError(t, err) @@ -92,8 +90,7 @@ func TestGetDeploymentWorkspaceAgentStats(t *testing.T) { TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 1, - SessionCountVSCode: 1, - }) + }, map[string]int64{"vscode": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ // Ensure this stat is newer! CreatedAt: insertTime, @@ -101,8 +98,7 @@ func TestGetDeploymentWorkspaceAgentStats(t *testing.T) { TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 2, - SessionCountVSCode: 1, - }) + }, map[string]int64{"vscode": 1}) stats, err := db.GetDeploymentWorkspaceAgentStats(ctx, dbtime.Now().Add(-time.Hour)) require.NoError(t, err) @@ -135,22 +131,17 @@ func TestGetDeploymentWorkspaceAgentUsageStats(t *testing.T) { TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 1, - // Should be ignored - SessionCountSSH: 4, - SessionCountVSCode: 3, - }) + }, map[string]int64{"ssh": 4, "vscode": 3}) // Should be ignored. dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime.Add(-time.Minute), - AgentID: agentID, - SessionCountVSCode: 1, - Usage: true, - }) + CreatedAt: insertTime.Add(-time.Minute), + AgentID: agentID, + Usage: true, + }, map[string]int64{"vscode": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime.Add(-time.Minute), - AgentID: agentID, - SessionCountReconnectingPTY: 1, - Usage: true, - }) + CreatedAt: insertTime.Add(-time.Minute), + AgentID: agentID, + Usage: true, + }, map[string]int64{"reconnecting_pty": 1}) // Latest stats dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ @@ -159,22 +150,17 @@ func TestGetDeploymentWorkspaceAgentUsageStats(t *testing.T) { TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 2, - // Should be ignored - SessionCountSSH: 3, - SessionCountVSCode: 1, - }) + }, map[string]int64{"ssh": 3, "vscode": 1}) // Should be ignored. dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agentID, - SessionCountVSCode: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agentID, + Usage: true, + }, map[string]int64{"vscode": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agentID, - SessionCountSSH: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agentID, + Usage: true, + }, map[string]int64{"ssh": 1}) stats, err := db.GetDeploymentWorkspaceAgentUsageStats(ctx, dbtime.Now().Add(-time.Hour)) require.NoError(t, err) @@ -189,6 +175,34 @@ func TestGetDeploymentWorkspaceAgentUsageStats(t *testing.T) { require.Equal(t, int64(0), stats.SessionCountJetBrains) }) + t.Run("ExcludesStatsBeforeCutoffInSameMinute", func(t *testing.T) { + t.Parallel() + + db, _ := dbtestutil.NewDB(t) + authz := rbac.NewAuthorizer(prometheus.NewRegistry()) + db = dbauthz.New(db, authz, slogtest.Make(t, &slogtest.Options{}), coderdtest.AccessControlStorePointer()) + ctx := context.Background() + agentID := uuid.New() + minute := dbtime.Now().Add(-2 * time.Minute).Truncate(time.Minute) + cutoff := minute.Add(30 * time.Second) + + dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ + CreatedAt: minute.Add(10 * time.Second), + AgentID: agentID, + Usage: true, + }, map[string]int64{"vscode": 4}) + dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ + CreatedAt: minute.Add(40 * time.Second), + AgentID: agentID, + Usage: true, + }, map[string]int64{"ssh": 1}) + + stats, err := db.GetDeploymentWorkspaceAgentUsageStats(ctx, cutoff) + require.NoError(t, err) + require.Zero(t, stats.SessionCountVSCode) + require.Equal(t, int64(1), stats.SessionCountSSH) + }) + t.Run("NoUsage", func(t *testing.T) { t.Parallel() @@ -206,10 +220,7 @@ func TestGetDeploymentWorkspaceAgentUsageStats(t *testing.T) { TxBytes: 3, RxBytes: 4, ConnectionMedianLatencyMS: 2, - // Should be ignored - SessionCountSSH: 3, - SessionCountVSCode: 1, - }) + }, map[string]int64{"ssh": 3, "vscode": 1}) // Should be ignored. stats, err := db.GetDeploymentWorkspaceAgentUsageStats(ctx, dbtime.Now().Add(-time.Hour)) require.NoError(t, err) @@ -709,6 +720,62 @@ func TestGetProvisionerDaemonsWithStatusByOrganization(t *testing.T) { }) } +func TestGetTemplateInsightsByTemplate(t *testing.T) { + t.Parallel() + + db, _ := dbtestutil.NewDB(t) + ctx := context.Background() + startTime := dbtime.Now().Add(-10 * time.Minute).Truncate(time.Minute) + endTime := startTime.Add(2 * time.Minute) + + insertStat := func(offset time.Duration, templateID, userID, workspaceID uuid.UUID, connections int64, counts map[string]int64) { + dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ + CreatedAt: startTime.Add(offset), + TemplateID: templateID, + UserID: userID, + WorkspaceID: workspaceID, + AgentID: uuid.New(), + ConnectionCount: connections, + }, counts) + } + + templateID := uuid.New() + userID := uuid.New() + workspaceID := uuid.New() + otherWorkspaceID := uuid.New() + + // Activity is deduplicated by user and minute. + insertStat(5*time.Second, templateID, userID, workspaceID, 1, map[string]int64{"ssh": 1, "vscode": 1}) + insertStat(20*time.Second, templateID, userID, workspaceID, 0, map[string]int64{"jetbrains": 1, "vscode": 1}) + insertStat(40*time.Second, templateID, userID, otherWorkspaceID, 0, map[string]int64{"reconnecting_pty": 1, "vscode": 1}) + insertStat(time.Minute+5*time.Second, templateID, userID, otherWorkspaceID, 1, map[string]int64{"ssh": 1, "vscode": 1}) + + // Unknown apps do not contribute activity or active users. + insertStat(10*time.Second, templateID, uuid.New(), uuid.New(), 1, map[string]int64{"unknown": 1}) + + // Unknown activity cannot supply a connection for known activity. + noKnownConnectionTemplateID := uuid.New() + noKnownConnectionUserID := uuid.New() + insertStat(15*time.Second, noKnownConnectionTemplateID, noKnownConnectionUserID, uuid.New(), 0, map[string]int64{"vscode": 1}) + insertStat(30*time.Second, noKnownConnectionTemplateID, noKnownConnectionUserID, uuid.New(), 1, map[string]int64{"unknown": 1}) + + insights, err := db.GetTemplateInsightsByTemplate(ctx, database.GetTemplateInsightsByTemplateParams{ + StartTime: startTime, + EndTime: endTime, + }) + require.NoError(t, err) + require.Equal(t, []database.GetTemplateInsightsByTemplateRow{ + { + TemplateID: templateID, + ActiveUsers: 1, + UsageVscodeSeconds: 120, + UsageJetbrainsSeconds: 60, + UsageReconnectingPtySeconds: 60, + UsageSshSeconds: 120, + }, + }, insights) +} + func TestGetWorkspaceAgentUsageStats(t *testing.T) { t.Parallel() @@ -741,19 +808,15 @@ func TestGetWorkspaceAgentUsageStats(t *testing.T) { TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 1, - // Should be ignored - SessionCountVSCode: 3, - SessionCountSSH: 1, - }) + }, map[string]int64{"vscode": 3, "ssh": 1}) // Should be ignored. dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime.Add(-time.Minute), - AgentID: agentID1, - WorkspaceID: workspaceID1, - TemplateID: templateID1, - UserID: userID1, - SessionCountVSCode: 1, - Usage: true, - }) + CreatedAt: insertTime.Add(-time.Minute), + AgentID: agentID1, + WorkspaceID: workspaceID1, + TemplateID: templateID1, + UserID: userID1, + Usage: true, + }, map[string]int64{"vscode": 1}) // Latest workspace 1 stats dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ @@ -765,28 +828,23 @@ func TestGetWorkspaceAgentUsageStats(t *testing.T) { TxBytes: 2, RxBytes: 2, ConnectionMedianLatencyMS: 1, - // Should be ignored - SessionCountVSCode: 3, - SessionCountSSH: 4, - }) + }, map[string]int64{"vscode": 3, "ssh": 4}) // Should be ignored. dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agentID1, - WorkspaceID: workspaceID1, - TemplateID: templateID1, - UserID: userID1, - SessionCountVSCode: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agentID1, + WorkspaceID: workspaceID1, + TemplateID: templateID1, + UserID: userID1, + Usage: true, + }, map[string]int64{"vscode": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agentID1, - WorkspaceID: workspaceID1, - TemplateID: templateID1, - UserID: userID1, - SessionCountJetBrains: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agentID1, + WorkspaceID: workspaceID1, + TemplateID: templateID1, + UserID: userID1, + Usage: true, + }, map[string]int64{"jetbrains": 1}) // Latest workspace 2 stats dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ @@ -808,28 +866,23 @@ func TestGetWorkspaceAgentUsageStats(t *testing.T) { TxBytes: 2, RxBytes: 3, ConnectionMedianLatencyMS: 1, - // Should be ignored - SessionCountVSCode: 3, - SessionCountSSH: 4, - }) + }, map[string]int64{"vscode": 3, "ssh": 4}) // Should be ignored. dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agentID2, - WorkspaceID: workspaceID2, - TemplateID: templateID2, - UserID: userID2, - SessionCountSSH: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agentID2, + WorkspaceID: workspaceID2, + TemplateID: templateID2, + UserID: userID2, + Usage: true, + }, map[string]int64{"ssh": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agentID2, - WorkspaceID: workspaceID2, - TemplateID: templateID2, - UserID: userID2, - SessionCountJetBrains: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agentID2, + WorkspaceID: workspaceID2, + TemplateID: templateID2, + UserID: userID2, + Usage: true, + }, map[string]int64{"jetbrains": 1}) reqTime := dbtime.Now().Add(-time.Hour) stats, err := db.GetWorkspaceAgentUsageStats(ctx, reqTime) @@ -872,10 +925,7 @@ func TestGetWorkspaceAgentUsageStats(t *testing.T) { TxBytes: 3, RxBytes: 4, ConnectionMedianLatencyMS: 2, - // Should be ignored - SessionCountSSH: 3, - SessionCountVSCode: 1, - }) + }, map[string]int64{"ssh": 3, "vscode": 1}) // Should be ignored. stats, err := db.GetWorkspaceAgentUsageStats(ctx, dbtime.Now().Add(-time.Hour)) require.NoError(t, err) @@ -1039,19 +1089,15 @@ func TestGetWorkspaceAgentUsageStatsAndLabels(t *testing.T) { TxBytes: 1, RxBytes: 1, ConnectionMedianLatencyMS: 1, - // Should be ignored - SessionCountVSCode: 3, - SessionCountSSH: 1, - }) + }, map[string]int64{"vscode": 3, "ssh": 1}) // Should be ignored. dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime.Add(-time.Minute), - AgentID: agent1.ID, - WorkspaceID: workspace1.ID, - TemplateID: template1.ID, - UserID: user1.ID, - SessionCountVSCode: 1, - Usage: true, - }) + CreatedAt: insertTime.Add(-time.Minute), + AgentID: agent1.ID, + WorkspaceID: workspace1.ID, + TemplateID: template1.ID, + UserID: user1.ID, + Usage: true, + }, map[string]int64{"vscode": 1}) // Latest workspace 1 stats dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ @@ -1063,28 +1109,23 @@ func TestGetWorkspaceAgentUsageStatsAndLabels(t *testing.T) { TxBytes: 2, RxBytes: 2, ConnectionMedianLatencyMS: 1, - // Should be ignored - SessionCountVSCode: 4, - SessionCountSSH: 3, - }) + }, map[string]int64{"vscode": 4, "ssh": 3}) // Should be ignored. dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agent1.ID, - WorkspaceID: workspace1.ID, - TemplateID: template1.ID, - UserID: user1.ID, - SessionCountJetBrains: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agent1.ID, + WorkspaceID: workspace1.ID, + TemplateID: template1.ID, + UserID: user1.ID, + Usage: true, + }, map[string]int64{"jetbrains": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agent1.ID, - WorkspaceID: workspace1.ID, - TemplateID: template1.ID, - UserID: user1.ID, - SessionCountReconnectingPTY: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agent1.ID, + WorkspaceID: workspace1.ID, + TemplateID: template1.ID, + UserID: user1.ID, + Usage: true, + }, map[string]int64{"reconnecting_pty": 1}) // Latest workspace 2 stats dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ @@ -1098,23 +1139,21 @@ func TestGetWorkspaceAgentUsageStatsAndLabels(t *testing.T) { ConnectionMedianLatencyMS: 1, }) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agent2.ID, - WorkspaceID: workspace2.ID, - TemplateID: template2.ID, - UserID: user2.ID, - SessionCountVSCode: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agent2.ID, + WorkspaceID: workspace2.ID, + TemplateID: template2.ID, + UserID: user2.ID, + Usage: true, + }, map[string]int64{"vscode": 1}) dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ - CreatedAt: insertTime, - AgentID: agent2.ID, - WorkspaceID: workspace2.ID, - TemplateID: template2.ID, - UserID: user2.ID, - SessionCountSSH: 1, - Usage: true, - }) + CreatedAt: insertTime, + AgentID: agent2.ID, + WorkspaceID: workspace2.ID, + TemplateID: template2.ID, + UserID: user2.ID, + Usage: true, + }, map[string]int64{"ssh": 1}) stats, err := db.GetWorkspaceAgentUsageStatsAndLabels(ctx, insertTime.Add(-time.Hour)) require.NoError(t, err) @@ -1180,10 +1219,7 @@ func TestGetWorkspaceAgentUsageStatsAndLabels(t *testing.T) { RxBytes: 4, TxBytes: 5, ConnectionMedianLatencyMS: 1, - // Should be ignored - SessionCountVSCode: 3, - SessionCountSSH: 1, - }) + }, map[string]int64{"vscode": 3, "ssh": 1}) // Should be ignored. stats, err := db.GetWorkspaceAgentUsageStatsAndLabels(ctx, insertTime.Add(-time.Hour)) require.NoError(t, err) diff --git a/coderd/database/queries.sql.go b/coderd/database/queries.sql.go index bb57a04bb3078..b5f61d95c0d5d 100644 --- a/coderd/database/queries.sql.go +++ b/coderd/database/queries.sql.go @@ -16012,39 +16012,41 @@ func (q *sqlQuerier) GetTemplateInsightsByInterval(ctx context.Context, arg GetT const getTemplateInsightsByTemplate = `-- name: GetTemplateInsightsByTemplate :many WITH - -- This CTE is used to truncate agent usage into minute buckets, then - -- flatten the users agent usage within the template so that usage in - -- multiple workspaces under one template is only counted once for - -- every minute (per user). + -- Deduplicate activity by template, user, and minute. + minute_activity AS ( + SELECT + template_id, + user_id, + date_trunc('minute', created_at) AS minute, + BOOL_OR(session_counts ? 'ssh') AS ssh, + BOOL_OR(session_counts ? 'reconnecting_pty') AS reconnecting_pty, + BOOL_OR(session_counts ? 'vscode') AS vscode, + BOOL_OR(session_counts ? 'jetbrains') AS jetbrains, + BOOL_OR(connection_count > 0) AS has_connection + FROM + workspace_agent_stats + WHERE + created_at >= $1::timestamptz + AND created_at < $2::timestamptz + AND session_counts ?| ARRAY['ssh', 'reconnecting_pty', 'vscode', 'jetbrains'] + GROUP BY + template_id, user_id, minute + ), insights AS ( SELECT template_id, user_id, - COUNT(DISTINCT CASE WHEN session_count_ssh > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS ssh_mins, - -- TODO(mafredri): Enable when we have the column. - -- COUNT(DISTINCT CASE WHEN session_count_sftp > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS sftp_mins, - COUNT(DISTINCT CASE WHEN session_count_reconnecting_pty > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS reconnecting_pty_mins, - COUNT(DISTINCT CASE WHEN session_count_vscode > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS vscode_mins, - COUNT(DISTINCT CASE WHEN session_count_jetbrains > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS jetbrains_mins, + COUNT(*) FILTER (WHERE ssh) AS ssh_mins, + COUNT(*) FILTER (WHERE reconnecting_pty) AS reconnecting_pty_mins, + COUNT(*) FILTER (WHERE vscode) AS vscode_mins, + COUNT(*) FILTER (WHERE jetbrains) AS jetbrains_mins, -- NOTE(mafredri): The agent stats are currently very unreliable, and -- sometimes the connections are missing, even during active sessions. -- Since we can't fully rely on this, we check for "any connection -- within this bucket". A better solution here would be preferable. - MAX(connection_count) > 0 AS has_connection + BOOL_OR(has_connection) AS has_connection FROM - workspace_agent_stats - WHERE - created_at >= $1::timestamptz - AND created_at < $2::timestamptz - -- Inclusion criteria to filter out empty results. - AND ( - session_count_ssh > 0 - -- TODO(mafredri): Enable when we have the column. - -- OR session_count_sftp > 0 - OR session_count_reconnecting_pty > 0 - OR session_count_vscode > 0 - OR session_count_jetbrains > 0 - ) + minute_activity GROUP BY template_id, user_id ) @@ -16654,27 +16656,11 @@ WITH template_id, user_id, -- Store each unique minute bucket for later merge between datasets. - array_agg( - DISTINCT CASE - WHEN - session_count_ssh > 0 - -- TODO(mafredri): Enable when we have the column. - -- OR session_count_sftp > 0 - OR session_count_reconnecting_pty > 0 - OR session_count_vscode > 0 - OR session_count_jetbrains > 0 - THEN - date_trunc('minute', created_at) - ELSE - NULL - END - ) AS minute_buckets, - COUNT(DISTINCT CASE WHEN session_count_ssh > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS ssh_mins, - -- TODO(mafredri): Enable when we have the column. - -- COUNT(DISTINCT CASE WHEN session_count_sftp > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS sftp_mins, - COUNT(DISTINCT CASE WHEN session_count_reconnecting_pty > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS reconnecting_pty_mins, - COUNT(DISTINCT CASE WHEN session_count_vscode > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS vscode_mins, - COUNT(DISTINCT CASE WHEN session_count_jetbrains > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS jetbrains_mins, + array_agg(DISTINCT date_trunc('minute', created_at)) AS minute_buckets, + COUNT(DISTINCT CASE WHEN session_counts ? 'ssh' THEN date_trunc('minute', created_at) ELSE NULL END) AS ssh_mins, + COUNT(DISTINCT CASE WHEN session_counts ? 'reconnecting_pty' THEN date_trunc('minute', created_at) ELSE NULL END) AS reconnecting_pty_mins, + COUNT(DISTINCT CASE WHEN session_counts ? 'vscode' THEN date_trunc('minute', created_at) ELSE NULL END) AS vscode_mins, + COUNT(DISTINCT CASE WHEN session_counts ? 'jetbrains' THEN date_trunc('minute', created_at) ELSE NULL END) AS jetbrains_mins, -- NOTE(mafredri): The agent stats are currently very unreliable, and -- sometimes the connections are missing, even during active sessions. -- Since we can't fully rely on this, we check for "any connection @@ -16687,15 +16673,7 @@ WITH -- AND created_at < @end_time::timestamptz created_at >= (SELECT t FROM latest_start) AND created_at < NOW() - -- Inclusion criteria to filter out empty results. - AND ( - session_count_ssh > 0 - -- TODO(mafredri): Enable when we have the column. - -- OR session_count_sftp > 0 - OR session_count_reconnecting_pty > 0 - OR session_count_vscode > 0 - OR session_count_jetbrains > 0 - ) + AND session_counts ?| ARRAY['ssh', 'reconnecting_pty', 'vscode', 'jetbrains'] GROUP BY time_bucket, template_id, user_id ), @@ -34839,30 +34817,27 @@ func (q *sqlQuerier) DeleteOldWorkspaceAgentStats(ctx context.Context) error { const getDeploymentWorkspaceAgentStats = `-- name: GetDeploymentWorkspaceAgentStats :one WITH stats AS ( - SELECT - agent_id, - created_at, - rx_bytes, - tx_bytes, - connection_median_latency_ms, - session_count_vscode, - session_count_ssh, - session_count_jetbrains, - session_count_reconnecting_pty, - ROW_NUMBER() OVER (PARTITION BY agent_id ORDER BY created_at DESC) AS rn - FROM workspace_agent_stats - WHERE created_at > $1 + SELECT + agent_id, + created_at, + rx_bytes, + tx_bytes, + connection_median_latency_ms, + session_counts, + ROW_NUMBER() OVER (PARTITION BY agent_id ORDER BY created_at DESC) AS rn + FROM workspace_agent_stats + WHERE created_at > $1 ) SELECT - coalesce(SUM(rx_bytes), 0)::bigint AS workspace_rx_bytes, - coalesce(SUM(tx_bytes), 0)::bigint AS workspace_tx_bytes, - -- The greater than 0 is to support legacy agents that don't report connection_median_latency_ms. - coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_50, - coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_95, - coalesce(SUM(session_count_vscode) FILTER (WHERE rn = 1), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh) FILTER (WHERE rn = 1), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains) FILTER (WHERE rn = 1), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty) FILTER (WHERE rn = 1), 0)::bigint AS session_count_reconnecting_pty + coalesce(SUM(rx_bytes), 0)::bigint AS workspace_rx_bytes, + coalesce(SUM(tx_bytes), 0)::bigint AS workspace_tx_bytes, + -- Positive latency values exclude legacy agents that do not report latency. + coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_50, + coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_95, + coalesce(SUM((session_counts ->> 'vscode')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_vscode, + coalesce(SUM((session_counts ->> 'ssh')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_ssh, + coalesce(SUM((session_counts ->> 'jetbrains')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((session_counts ->> 'reconnecting_pty')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_reconnecting_pty FROM stats ` @@ -34904,46 +34879,36 @@ WITH agent_stats AS ( -- The greater than 0 is to support legacy agents that don't report connection_median_latency_ms. WHERE workspace_agent_stats.created_at > $1 AND connection_median_latency_ms > 0 ), -minute_buckets AS ( +latest_minutes AS ( SELECT agent_id, - date_trunc('minute', created_at) AS minute_bucket, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty + MAX(date_trunc('minute', created_at)) AS minute_bucket FROM workspace_agent_stats WHERE created_at >= $1 - AND created_at < date_trunc('minute', now()) -- Exclude current partial minute - AND usage = true + -- Exclude the current partial minute. + AND created_at < date_trunc('minute', now()) + AND usage GROUP BY - agent_id, - minute_bucket -), -latest_buckets AS ( - SELECT DISTINCT ON (agent_id) - agent_id, - minute_bucket, - session_count_vscode, - session_count_jetbrains, - session_count_reconnecting_pty, - session_count_ssh - FROM - minute_buckets - ORDER BY - agent_id, - minute_bucket DESC + agent_id ), latest_agent_stats AS ( - SELECT - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty - FROM - latest_buckets + SELECT + coalesce(SUM((stats.session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((stats.session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((stats.session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((stats.session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty + FROM + latest_minutes + JOIN + workspace_agent_stats AS stats + ON + stats.agent_id = latest_minutes.agent_id + AND stats.created_at >= $1 + AND stats.created_at >= latest_minutes.minute_bucket + AND stats.created_at < latest_minutes.minute_bucket + '1 minute'::interval + AND stats.usage ) SELECT workspace_rx_bytes, workspace_tx_bytes, workspace_connection_latency_50, workspace_connection_latency_95, session_count_vscode, session_count_ssh, session_count_jetbrains, session_count_reconnecting_pty FROM agent_stats, latest_agent_stats ` @@ -34994,14 +34959,16 @@ WITH agent_stats AS ( ), latest_agent_stats AS ( SELECT a.agent_id, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty + coalesce(SUM((a.session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((a.session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((a.session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((a.session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty FROM ( - SELECT id, created_at, user_id, agent_id, workspace_id, template_id, connections_by_proto, connection_count, rx_packets, rx_bytes, tx_packets, tx_bytes, connection_median_latency_ms, session_count_vscode, session_count_jetbrains, session_count_reconnecting_pty, session_count_ssh, usage, ROW_NUMBER() OVER(PARTITION BY agent_id ORDER BY created_at DESC) AS rn + SELECT id, created_at, user_id, agent_id, workspace_id, template_id, connections_by_proto, connection_count, rx_packets, rx_bytes, tx_packets, tx_bytes, connection_median_latency_ms, usage, session_counts, ROW_NUMBER() OVER(PARTITION BY agent_id ORDER BY created_at DESC) AS rn FROM workspace_agent_stats WHERE created_at > $1 - ) AS a WHERE a.rn = 1 GROUP BY a.user_id, a.agent_id, a.workspace_id, a.template_id + ) AS a + WHERE a.rn = 1 + GROUP BY a.user_id, a.agent_id, a.workspace_id, a.template_id ) SELECT user_id, agent_stats.agent_id, workspace_id, template_id, aggregated_from, workspace_rx_bytes, workspace_tx_bytes, workspace_connection_latency_50, workspace_connection_latency_95, latest_agent_stats.agent_id, session_count_vscode, session_count_ssh, session_count_jetbrains, session_count_reconnecting_pty FROM agent_stats JOIN latest_agent_stats ON agent_stats.agent_id = latest_agent_stats.agent_id ` @@ -35075,14 +35042,14 @@ WITH agent_stats AS ( ), latest_agent_stats AS ( SELECT a.agent_id, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty, - coalesce(SUM(connection_count), 0)::bigint AS connection_count, - coalesce(MAX(connection_median_latency_ms), 0)::float AS connection_median_latency_ms + coalesce(SUM((a.session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((a.session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((a.session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((a.session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty, + coalesce(SUM(a.connection_count), 0)::bigint AS connection_count, + coalesce(MAX(a.connection_median_latency_ms), 0)::float AS connection_median_latency_ms FROM ( - SELECT id, created_at, user_id, agent_id, workspace_id, template_id, connections_by_proto, connection_count, rx_packets, rx_bytes, tx_packets, tx_bytes, connection_median_latency_ms, session_count_vscode, session_count_jetbrains, session_count_reconnecting_pty, session_count_ssh, usage, ROW_NUMBER() OVER(PARTITION BY agent_id ORDER BY created_at DESC) AS rn + SELECT id, created_at, user_id, agent_id, workspace_id, template_id, connections_by_proto, connection_count, rx_packets, rx_bytes, tx_packets, tx_bytes, connection_median_latency_ms, usage, session_counts, ROW_NUMBER() OVER(PARTITION BY agent_id ORDER BY created_at DESC) AS rn FROM workspace_agent_stats -- The greater than 0 is to support legacy agents that don't report connection_median_latency_ms. WHERE created_at > $1 AND connection_median_latency_ms > 0 @@ -35164,72 +35131,60 @@ func (q *sqlQuerier) GetWorkspaceAgentStatsAndLabels(ctx context.Context, create } const getWorkspaceAgentUsageStats = `-- name: GetWorkspaceAgentUsageStats :many -WITH agent_stats AS ( +WITH stats AS ( SELECT user_id, agent_id, workspace_id, template_id, - MIN(created_at)::timestamptz AS aggregated_from, - coalesce(SUM(rx_bytes), 0)::bigint AS workspace_rx_bytes, - coalesce(SUM(tx_bytes), 0)::bigint AS workspace_tx_bytes, - coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms)), -1)::FLOAT AS workspace_connection_latency_50, - coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms)), -1)::FLOAT AS workspace_connection_latency_95 + created_at, + rx_bytes, + tx_bytes, + connection_median_latency_ms, + usage, + session_counts, + MAX(date_trunc('minute', created_at)) FILTER ( + WHERE usage AND created_at < date_trunc('minute', now()) + ) OVER (PARTITION BY agent_id) AS latest_usage_minute FROM workspace_agent_stats - -- The greater than 0 is to support legacy agents that don't report connection_median_latency_ms. - WHERE workspace_agent_stats.created_at > $1 AND connection_median_latency_ms > 0 - GROUP BY user_id, agent_id, workspace_id, template_id -), -minute_buckets AS ( - SELECT - agent_id, - date_trunc('minute', created_at) AS minute_bucket, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty - FROM - workspace_agent_stats - WHERE - created_at >= $1 - AND created_at < date_trunc('minute', now()) -- Exclude current partial minute - AND usage = true - GROUP BY - agent_id, - minute_bucket, - user_id, - agent_id, - workspace_id, - template_id -), -latest_buckets AS ( - SELECT DISTINCT ON (agent_id) - agent_id, - session_count_vscode, - session_count_ssh, - session_count_jetbrains, - session_count_reconnecting_pty - FROM - minute_buckets - ORDER BY - agent_id, - minute_bucket DESC + WHERE created_at >= $1 ) -SELECT user_id, -agent_stats.agent_id, -workspace_id, -template_id, -aggregated_from, -workspace_rx_bytes, -workspace_tx_bytes, -workspace_connection_latency_50, -workspace_connection_latency_95, -coalesce(latest_buckets.agent_id,agent_stats.agent_id) AS agent_id, -coalesce(session_count_vscode, 0)::bigint AS session_count_vscode, -coalesce(session_count_ssh, 0)::bigint AS session_count_ssh, -coalesce(session_count_jetbrains, 0)::bigint AS session_count_jetbrains, -coalesce(session_count_reconnecting_pty, 0)::bigint AS session_count_reconnecting_pty -FROM agent_stats LEFT JOIN latest_buckets ON agent_stats.agent_id = latest_buckets.agent_id +SELECT + user_id, + agent_id, + workspace_id, + template_id, + MIN(created_at) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + )::timestamptz AS aggregated_from, + coalesce(SUM(rx_bytes) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + ), 0)::bigint AS workspace_rx_bytes, + coalesce(SUM(tx_bytes) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + ), 0)::bigint AS workspace_tx_bytes, + coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + )), -1)::FLOAT AS workspace_connection_latency_50, + coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + )), -1)::FLOAT AS workspace_connection_latency_95, + agent_id, + coalesce(SUM((session_counts ->> 'vscode')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_vscode, + coalesce(SUM((session_counts ->> 'ssh')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_ssh, + coalesce(SUM((session_counts ->> 'jetbrains')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((session_counts ->> 'reconnecting_pty')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_reconnecting_pty +FROM stats +GROUP BY user_id, agent_id, workspace_id, template_id +HAVING BOOL_OR(created_at > $1 AND connection_median_latency_ms > 0) ` type GetWorkspaceAgentUsageStatsRow struct { @@ -35249,7 +35204,6 @@ type GetWorkspaceAgentUsageStatsRow struct { SessionCountReconnectingPTY int64 `db:"session_count_reconnecting_pty" json:"session_count_reconnecting_pty"` } -// `minute_buckets` could return 0 rows if there are no usage stats since `created_at`. func (q *sqlQuerier) GetWorkspaceAgentUsageStats(ctx context.Context, createdAt time.Time) ([]GetWorkspaceAgentUsageStatsRow, error) { rows, err := q.db.QueryContext(ctx, getWorkspaceAgentUsageStats, createdAt) if err != nil { @@ -35304,15 +35258,15 @@ WITH agent_stats AS ( ), latest_agent_stats AS ( SELECT agent_id, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty, + coalesce(SUM((session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty, coalesce(SUM(connection_count), 0)::bigint AS connection_count FROM workspace_agent_stats -- We only want the latest stats, but those stats might be -- spread across multiple rows. - WHERE usage = true AND created_at > now() - '1 minute'::interval + WHERE usage AND created_at > now() - '1 minute'::interval GROUP BY user_id, agent_id, workspace_id ) SELECT @@ -35407,10 +35361,7 @@ INSERT INTO rx_bytes, tx_packets, tx_bytes, - session_count_vscode, - session_count_jetbrains, - session_count_reconnecting_pty, - session_count_ssh, + session_counts, connection_median_latency_ms, usage ) @@ -35427,33 +35378,27 @@ SELECT unnest($10 :: bigint[]) AS rx_bytes, unnest($11 :: bigint[]) AS tx_packets, unnest($12 :: bigint[]) AS tx_bytes, - unnest($13 :: bigint[]) AS session_count_vscode, - unnest($14 :: bigint[]) AS session_count_jetbrains, - unnest($15 :: bigint[]) AS session_count_reconnecting_pty, - unnest($16 :: bigint[]) AS session_count_ssh, - unnest($17 :: double precision[]) AS connection_median_latency_ms, - unnest($18 :: boolean[]) AS usage + jsonb_array_elements($13 :: jsonb) AS session_counts, + unnest($14 :: double precision[]) AS connection_median_latency_ms, + unnest($15 :: boolean[]) AS usage ` type InsertWorkspaceAgentStatsParams struct { - ID []uuid.UUID `db:"id" json:"id"` - CreatedAt []time.Time `db:"created_at" json:"created_at"` - UserID []uuid.UUID `db:"user_id" json:"user_id"` - WorkspaceID []uuid.UUID `db:"workspace_id" json:"workspace_id"` - TemplateID []uuid.UUID `db:"template_id" json:"template_id"` - AgentID []uuid.UUID `db:"agent_id" json:"agent_id"` - ConnectionsByProto json.RawMessage `db:"connections_by_proto" json:"connections_by_proto"` - ConnectionCount []int64 `db:"connection_count" json:"connection_count"` - RxPackets []int64 `db:"rx_packets" json:"rx_packets"` - RxBytes []int64 `db:"rx_bytes" json:"rx_bytes"` - TxPackets []int64 `db:"tx_packets" json:"tx_packets"` - TxBytes []int64 `db:"tx_bytes" json:"tx_bytes"` - SessionCountVSCode []int64 `db:"session_count_vscode" json:"session_count_vscode"` - SessionCountJetBrains []int64 `db:"session_count_jetbrains" json:"session_count_jetbrains"` - SessionCountReconnectingPTY []int64 `db:"session_count_reconnecting_pty" json:"session_count_reconnecting_pty"` - SessionCountSSH []int64 `db:"session_count_ssh" json:"session_count_ssh"` - ConnectionMedianLatencyMS []float64 `db:"connection_median_latency_ms" json:"connection_median_latency_ms"` - Usage []bool `db:"usage" json:"usage"` + ID []uuid.UUID `db:"id" json:"id"` + CreatedAt []time.Time `db:"created_at" json:"created_at"` + UserID []uuid.UUID `db:"user_id" json:"user_id"` + WorkspaceID []uuid.UUID `db:"workspace_id" json:"workspace_id"` + TemplateID []uuid.UUID `db:"template_id" json:"template_id"` + AgentID []uuid.UUID `db:"agent_id" json:"agent_id"` + ConnectionsByProto json.RawMessage `db:"connections_by_proto" json:"connections_by_proto"` + ConnectionCount []int64 `db:"connection_count" json:"connection_count"` + RxPackets []int64 `db:"rx_packets" json:"rx_packets"` + RxBytes []int64 `db:"rx_bytes" json:"rx_bytes"` + TxPackets []int64 `db:"tx_packets" json:"tx_packets"` + TxBytes []int64 `db:"tx_bytes" json:"tx_bytes"` + SessionCounts json.RawMessage `db:"session_counts" json:"session_counts"` + ConnectionMedianLatencyMS []float64 `db:"connection_median_latency_ms" json:"connection_median_latency_ms"` + Usage []bool `db:"usage" json:"usage"` } func (q *sqlQuerier) InsertWorkspaceAgentStats(ctx context.Context, arg InsertWorkspaceAgentStatsParams) error { @@ -35470,10 +35415,7 @@ func (q *sqlQuerier) InsertWorkspaceAgentStats(ctx context.Context, arg InsertWo pq.Array(arg.RxBytes), pq.Array(arg.TxPackets), pq.Array(arg.TxBytes), - pq.Array(arg.SessionCountVSCode), - pq.Array(arg.SessionCountJetBrains), - pq.Array(arg.SessionCountReconnectingPTY), - pq.Array(arg.SessionCountSSH), + arg.SessionCounts, pq.Array(arg.ConnectionMedianLatencyMS), pq.Array(arg.Usage), ) diff --git a/coderd/database/queries/insights.sql b/coderd/database/queries/insights.sql index b589ce4e9a6fe..6e61b95bf0806 100644 --- a/coderd/database/queries/insights.sql +++ b/coderd/database/queries/insights.sql @@ -148,39 +148,41 @@ FROM -- GetTemplateInsightsByTemplate is used for Prometheus metrics. Keep -- in sync with GetTemplateInsights and UpsertTemplateUsageStats. WITH - -- This CTE is used to truncate agent usage into minute buckets, then - -- flatten the users agent usage within the template so that usage in - -- multiple workspaces under one template is only counted once for - -- every minute (per user). + -- Deduplicate activity by template, user, and minute. + minute_activity AS ( + SELECT + template_id, + user_id, + date_trunc('minute', created_at) AS minute, + BOOL_OR(session_counts ? 'ssh') AS ssh, + BOOL_OR(session_counts ? 'reconnecting_pty') AS reconnecting_pty, + BOOL_OR(session_counts ? 'vscode') AS vscode, + BOOL_OR(session_counts ? 'jetbrains') AS jetbrains, + BOOL_OR(connection_count > 0) AS has_connection + FROM + workspace_agent_stats + WHERE + created_at >= @start_time::timestamptz + AND created_at < @end_time::timestamptz + AND session_counts ?| ARRAY['ssh', 'reconnecting_pty', 'vscode', 'jetbrains'] + GROUP BY + template_id, user_id, minute + ), insights AS ( SELECT template_id, user_id, - COUNT(DISTINCT CASE WHEN session_count_ssh > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS ssh_mins, - -- TODO(mafredri): Enable when we have the column. - -- COUNT(DISTINCT CASE WHEN session_count_sftp > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS sftp_mins, - COUNT(DISTINCT CASE WHEN session_count_reconnecting_pty > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS reconnecting_pty_mins, - COUNT(DISTINCT CASE WHEN session_count_vscode > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS vscode_mins, - COUNT(DISTINCT CASE WHEN session_count_jetbrains > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS jetbrains_mins, + COUNT(*) FILTER (WHERE ssh) AS ssh_mins, + COUNT(*) FILTER (WHERE reconnecting_pty) AS reconnecting_pty_mins, + COUNT(*) FILTER (WHERE vscode) AS vscode_mins, + COUNT(*) FILTER (WHERE jetbrains) AS jetbrains_mins, -- NOTE(mafredri): The agent stats are currently very unreliable, and -- sometimes the connections are missing, even during active sessions. -- Since we can't fully rely on this, we check for "any connection -- within this bucket". A better solution here would be preferable. - MAX(connection_count) > 0 AS has_connection + BOOL_OR(has_connection) AS has_connection FROM - workspace_agent_stats - WHERE - created_at >= @start_time::timestamptz - AND created_at < @end_time::timestamptz - -- Inclusion criteria to filter out empty results. - AND ( - session_count_ssh > 0 - -- TODO(mafredri): Enable when we have the column. - -- OR session_count_sftp > 0 - OR session_count_reconnecting_pty > 0 - OR session_count_vscode > 0 - OR session_count_jetbrains > 0 - ) + minute_activity GROUP BY template_id, user_id ) @@ -559,27 +561,11 @@ WITH template_id, user_id, -- Store each unique minute bucket for later merge between datasets. - array_agg( - DISTINCT CASE - WHEN - session_count_ssh > 0 - -- TODO(mafredri): Enable when we have the column. - -- OR session_count_sftp > 0 - OR session_count_reconnecting_pty > 0 - OR session_count_vscode > 0 - OR session_count_jetbrains > 0 - THEN - date_trunc('minute', created_at) - ELSE - NULL - END - ) AS minute_buckets, - COUNT(DISTINCT CASE WHEN session_count_ssh > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS ssh_mins, - -- TODO(mafredri): Enable when we have the column. - -- COUNT(DISTINCT CASE WHEN session_count_sftp > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS sftp_mins, - COUNT(DISTINCT CASE WHEN session_count_reconnecting_pty > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS reconnecting_pty_mins, - COUNT(DISTINCT CASE WHEN session_count_vscode > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS vscode_mins, - COUNT(DISTINCT CASE WHEN session_count_jetbrains > 0 THEN date_trunc('minute', created_at) ELSE NULL END) AS jetbrains_mins, + array_agg(DISTINCT date_trunc('minute', created_at)) AS minute_buckets, + COUNT(DISTINCT CASE WHEN session_counts ? 'ssh' THEN date_trunc('minute', created_at) ELSE NULL END) AS ssh_mins, + COUNT(DISTINCT CASE WHEN session_counts ? 'reconnecting_pty' THEN date_trunc('minute', created_at) ELSE NULL END) AS reconnecting_pty_mins, + COUNT(DISTINCT CASE WHEN session_counts ? 'vscode' THEN date_trunc('minute', created_at) ELSE NULL END) AS vscode_mins, + COUNT(DISTINCT CASE WHEN session_counts ? 'jetbrains' THEN date_trunc('minute', created_at) ELSE NULL END) AS jetbrains_mins, -- NOTE(mafredri): The agent stats are currently very unreliable, and -- sometimes the connections are missing, even during active sessions. -- Since we can't fully rely on this, we check for "any connection @@ -592,15 +578,7 @@ WITH -- AND created_at < @end_time::timestamptz created_at >= (SELECT t FROM latest_start) AND created_at < NOW() - -- Inclusion criteria to filter out empty results. - AND ( - session_count_ssh > 0 - -- TODO(mafredri): Enable when we have the column. - -- OR session_count_sftp > 0 - OR session_count_reconnecting_pty > 0 - OR session_count_vscode > 0 - OR session_count_jetbrains > 0 - ) + AND session_counts ?| ARRAY['ssh', 'reconnecting_pty', 'vscode', 'jetbrains'] GROUP BY time_bucket, template_id, user_id ), diff --git a/coderd/database/queries/workspaceagentstats.sql b/coderd/database/queries/workspaceagentstats.sql index 28c17d8271e8d..a0c6298f5cbe7 100644 --- a/coderd/database/queries/workspaceagentstats.sql +++ b/coderd/database/queries/workspaceagentstats.sql @@ -13,10 +13,7 @@ INSERT INTO rx_bytes, tx_packets, tx_bytes, - session_count_vscode, - session_count_jetbrains, - session_count_reconnecting_pty, - session_count_ssh, + session_counts, connection_median_latency_ms, usage ) @@ -33,10 +30,7 @@ SELECT unnest(@rx_bytes :: bigint[]) AS rx_bytes, unnest(@tx_packets :: bigint[]) AS tx_packets, unnest(@tx_bytes :: bigint[]) AS tx_bytes, - unnest(@session_count_vscode :: bigint[]) AS session_count_vscode, - unnest(@session_count_jetbrains :: bigint[]) AS session_count_jetbrains, - unnest(@session_count_reconnecting_pty :: bigint[]) AS session_count_reconnecting_pty, - unnest(@session_count_ssh :: bigint[]) AS session_count_ssh, + jsonb_array_elements(@session_counts :: jsonb) AS session_counts, unnest(@connection_median_latency_ms :: double precision[]) AS connection_median_latency_ms, unnest(@usage :: boolean[]) AS usage; @@ -73,30 +67,27 @@ WHERE -- name: GetDeploymentWorkspaceAgentStats :one WITH stats AS ( - SELECT - agent_id, - created_at, - rx_bytes, - tx_bytes, - connection_median_latency_ms, - session_count_vscode, - session_count_ssh, - session_count_jetbrains, - session_count_reconnecting_pty, - ROW_NUMBER() OVER (PARTITION BY agent_id ORDER BY created_at DESC) AS rn - FROM workspace_agent_stats - WHERE created_at > $1 + SELECT + agent_id, + created_at, + rx_bytes, + tx_bytes, + connection_median_latency_ms, + session_counts, + ROW_NUMBER() OVER (PARTITION BY agent_id ORDER BY created_at DESC) AS rn + FROM workspace_agent_stats + WHERE created_at > $1 ) SELECT - coalesce(SUM(rx_bytes), 0)::bigint AS workspace_rx_bytes, - coalesce(SUM(tx_bytes), 0)::bigint AS workspace_tx_bytes, - -- The greater than 0 is to support legacy agents that don't report connection_median_latency_ms. - coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_50, - coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_95, - coalesce(SUM(session_count_vscode) FILTER (WHERE rn = 1), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh) FILTER (WHERE rn = 1), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains) FILTER (WHERE rn = 1), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty) FILTER (WHERE rn = 1), 0)::bigint AS session_count_reconnecting_pty + coalesce(SUM(rx_bytes), 0)::bigint AS workspace_rx_bytes, + coalesce(SUM(tx_bytes), 0)::bigint AS workspace_tx_bytes, + -- Positive latency values exclude legacy agents that do not report latency. + coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_50, + coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER (WHERE connection_median_latency_ms > 0)), -1)::FLOAT AS workspace_connection_latency_95, + coalesce(SUM((session_counts ->> 'vscode')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_vscode, + coalesce(SUM((session_counts ->> 'ssh')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_ssh, + coalesce(SUM((session_counts ->> 'jetbrains')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((session_counts ->> 'reconnecting_pty')::bigint) FILTER (WHERE rn = 1), 0)::bigint AS session_count_reconnecting_pty FROM stats; -- name: GetDeploymentWorkspaceAgentUsageStats :one @@ -110,46 +101,36 @@ WITH agent_stats AS ( -- The greater than 0 is to support legacy agents that don't report connection_median_latency_ms. WHERE workspace_agent_stats.created_at > $1 AND connection_median_latency_ms > 0 ), -minute_buckets AS ( +latest_minutes AS ( SELECT agent_id, - date_trunc('minute', created_at) AS minute_bucket, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty + MAX(date_trunc('minute', created_at)) AS minute_bucket FROM workspace_agent_stats WHERE created_at >= $1 - AND created_at < date_trunc('minute', now()) -- Exclude current partial minute - AND usage = true + -- Exclude the current partial minute. + AND created_at < date_trunc('minute', now()) + AND usage GROUP BY - agent_id, - minute_bucket -), -latest_buckets AS ( - SELECT DISTINCT ON (agent_id) - agent_id, - minute_bucket, - session_count_vscode, - session_count_jetbrains, - session_count_reconnecting_pty, - session_count_ssh - FROM - minute_buckets - ORDER BY - agent_id, - minute_bucket DESC + agent_id ), latest_agent_stats AS ( - SELECT - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty - FROM - latest_buckets + SELECT + coalesce(SUM((stats.session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((stats.session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((stats.session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((stats.session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty + FROM + latest_minutes + JOIN + workspace_agent_stats AS stats + ON + stats.agent_id = latest_minutes.agent_id + AND stats.created_at >= $1 + AND stats.created_at >= latest_minutes.minute_bucket + AND stats.created_at < latest_minutes.minute_bucket + '1 minute'::interval + AND stats.usage ) SELECT * FROM agent_stats, latest_agent_stats; @@ -172,85 +153,74 @@ WITH agent_stats AS ( ), latest_agent_stats AS ( SELECT a.agent_id, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty + coalesce(SUM((a.session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((a.session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((a.session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((a.session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty FROM ( SELECT *, ROW_NUMBER() OVER(PARTITION BY agent_id ORDER BY created_at DESC) AS rn FROM workspace_agent_stats WHERE created_at > $1 - ) AS a WHERE a.rn = 1 GROUP BY a.user_id, a.agent_id, a.workspace_id, a.template_id + ) AS a + WHERE a.rn = 1 + GROUP BY a.user_id, a.agent_id, a.workspace_id, a.template_id ) SELECT * FROM agent_stats JOIN latest_agent_stats ON agent_stats.agent_id = latest_agent_stats.agent_id; -- name: GetWorkspaceAgentUsageStats :many -WITH agent_stats AS ( +WITH stats AS ( SELECT user_id, agent_id, workspace_id, template_id, - MIN(created_at)::timestamptz AS aggregated_from, - coalesce(SUM(rx_bytes), 0)::bigint AS workspace_rx_bytes, - coalesce(SUM(tx_bytes), 0)::bigint AS workspace_tx_bytes, - coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms)), -1)::FLOAT AS workspace_connection_latency_50, - coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms)), -1)::FLOAT AS workspace_connection_latency_95 + created_at, + rx_bytes, + tx_bytes, + connection_median_latency_ms, + usage, + session_counts, + MAX(date_trunc('minute', created_at)) FILTER ( + WHERE usage AND created_at < date_trunc('minute', now()) + ) OVER (PARTITION BY agent_id) AS latest_usage_minute FROM workspace_agent_stats - -- The greater than 0 is to support legacy agents that don't report connection_median_latency_ms. - WHERE workspace_agent_stats.created_at > $1 AND connection_median_latency_ms > 0 - GROUP BY user_id, agent_id, workspace_id, template_id -), -minute_buckets AS ( - SELECT - agent_id, - date_trunc('minute', created_at) AS minute_bucket, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty - FROM - workspace_agent_stats - WHERE - created_at >= $1 - AND created_at < date_trunc('minute', now()) -- Exclude current partial minute - AND usage = true - GROUP BY - agent_id, - minute_bucket, - user_id, - agent_id, - workspace_id, - template_id -), -latest_buckets AS ( - SELECT DISTINCT ON (agent_id) - agent_id, - session_count_vscode, - session_count_ssh, - session_count_jetbrains, - session_count_reconnecting_pty - FROM - minute_buckets - ORDER BY - agent_id, - minute_bucket DESC + WHERE created_at >= $1 ) -SELECT user_id, -agent_stats.agent_id, -workspace_id, -template_id, -aggregated_from, -workspace_rx_bytes, -workspace_tx_bytes, -workspace_connection_latency_50, -workspace_connection_latency_95, --- `minute_buckets` could return 0 rows if there are no usage stats since `created_at`. -coalesce(latest_buckets.agent_id,agent_stats.agent_id) AS agent_id, -coalesce(session_count_vscode, 0)::bigint AS session_count_vscode, -coalesce(session_count_ssh, 0)::bigint AS session_count_ssh, -coalesce(session_count_jetbrains, 0)::bigint AS session_count_jetbrains, -coalesce(session_count_reconnecting_pty, 0)::bigint AS session_count_reconnecting_pty -FROM agent_stats LEFT JOIN latest_buckets ON agent_stats.agent_id = latest_buckets.agent_id; +SELECT + user_id, + agent_id, + workspace_id, + template_id, + MIN(created_at) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + )::timestamptz AS aggregated_from, + coalesce(SUM(rx_bytes) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + ), 0)::bigint AS workspace_rx_bytes, + coalesce(SUM(tx_bytes) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + ), 0)::bigint AS workspace_tx_bytes, + coalesce((PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + )), -1)::FLOAT AS workspace_connection_latency_50, + coalesce((PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY connection_median_latency_ms) FILTER ( + WHERE created_at > $1 AND connection_median_latency_ms > 0 + )), -1)::FLOAT AS workspace_connection_latency_95, + agent_id, + coalesce(SUM((session_counts ->> 'vscode')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_vscode, + coalesce(SUM((session_counts ->> 'ssh')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_ssh, + coalesce(SUM((session_counts ->> 'jetbrains')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((session_counts ->> 'reconnecting_pty')::bigint) FILTER ( + WHERE usage AND date_trunc('minute', created_at) = latest_usage_minute + ), 0)::bigint AS session_count_reconnecting_pty +FROM stats +GROUP BY user_id, agent_id, workspace_id, template_id +HAVING BOOL_OR(created_at > $1 AND connection_median_latency_ms > 0); -- name: GetWorkspaceAgentStatsAndLabels :many WITH agent_stats AS ( @@ -266,12 +236,12 @@ WITH agent_stats AS ( ), latest_agent_stats AS ( SELECT a.agent_id, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty, - coalesce(SUM(connection_count), 0)::bigint AS connection_count, - coalesce(MAX(connection_median_latency_ms), 0)::float AS connection_median_latency_ms + coalesce(SUM((a.session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((a.session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((a.session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((a.session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty, + coalesce(SUM(a.connection_count), 0)::bigint AS connection_count, + coalesce(MAX(a.connection_median_latency_ms), 0)::float AS connection_median_latency_ms FROM ( SELECT *, ROW_NUMBER() OVER(PARTITION BY agent_id ORDER BY created_at DESC) AS rn FROM workspace_agent_stats @@ -320,15 +290,15 @@ WITH agent_stats AS ( ), latest_agent_stats AS ( SELECT agent_id, - coalesce(SUM(session_count_vscode), 0)::bigint AS session_count_vscode, - coalesce(SUM(session_count_ssh), 0)::bigint AS session_count_ssh, - coalesce(SUM(session_count_jetbrains), 0)::bigint AS session_count_jetbrains, - coalesce(SUM(session_count_reconnecting_pty), 0)::bigint AS session_count_reconnecting_pty, + coalesce(SUM((session_counts ->> 'vscode')::bigint), 0)::bigint AS session_count_vscode, + coalesce(SUM((session_counts ->> 'ssh')::bigint), 0)::bigint AS session_count_ssh, + coalesce(SUM((session_counts ->> 'jetbrains')::bigint), 0)::bigint AS session_count_jetbrains, + coalesce(SUM((session_counts ->> 'reconnecting_pty')::bigint), 0)::bigint AS session_count_reconnecting_pty, coalesce(SUM(connection_count), 0)::bigint AS connection_count FROM workspace_agent_stats -- We only want the latest stats, but those stats might be -- spread across multiple rows. - WHERE usage = true AND created_at > now() - '1 minute'::interval + WHERE usage AND created_at > now() - '1 minute'::interval GROUP BY user_id, agent_id, workspace_id ) SELECT diff --git a/coderd/metricscache/metricscache_test.go b/coderd/metricscache/metricscache_test.go index f730dcc240058..6a7bf5b43a33f 100644 --- a/coderd/metricscache/metricscache_test.go +++ b/coderd/metricscache/metricscache_test.go @@ -3,7 +3,6 @@ package metricscache_test import ( "context" "database/sql" - "encoding/json" "sync/atomic" "testing" "time" @@ -308,28 +307,13 @@ func TestCache_DeploymentStats(t *testing.T) { DeploymentStats: time.Minute, }, false) - err := db.InsertWorkspaceAgentStats(context.Background(), database.InsertWorkspaceAgentStatsParams{ - ID: []uuid.UUID{uuid.New()}, - CreatedAt: []time.Time{clock.Now()}, - WorkspaceID: []uuid.UUID{uuid.New()}, - UserID: []uuid.UUID{uuid.New()}, - TemplateID: []uuid.UUID{uuid.New()}, - AgentID: []uuid.UUID{uuid.New()}, - ConnectionsByProto: json.RawMessage(`[{}]`), - - RxPackets: []int64{0}, - RxBytes: []int64{1}, - TxPackets: []int64{0}, - TxBytes: []int64{1}, - ConnectionCount: []int64{1}, - SessionCountVSCode: []int64{1}, - SessionCountJetBrains: []int64{0}, - SessionCountReconnectingPTY: []int64{0}, - SessionCountSSH: []int64{0}, - ConnectionMedianLatencyMS: []float64{10}, - Usage: []bool{false}, - }) - require.NoError(t, err) + dbgen.WorkspaceAgentStat(t, db, database.WorkspaceAgentStat{ + CreatedAt: clock.Now(), + RxBytes: 1, + TxBytes: 1, + ConnectionCount: 1, + ConnectionMedianLatencyMS: 10, + }, map[string]int64{"vscode": 1}) // Wait for both ticker functions to be created (template build times and deployment stats) tickerTrap.MustWait(ctx).MustRelease(ctx) diff --git a/coderd/workspacestats/batcher.go b/coderd/workspacestats/batcher.go index 847ef562fbb1c..b2c52d78456f7 100644 --- a/coderd/workspacestats/batcher.go +++ b/coderd/workspacestats/batcher.go @@ -17,6 +17,7 @@ import ( "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/idemetadata" ) const ( @@ -37,9 +38,9 @@ type DBBatcher struct { mu sync.Mutex // TODO: make this a buffered chan instead? buf *database.InsertWorkspaceAgentStatsParams - // NOTE: we batch this separately as it's a jsonb field and - // pq.Array + unnest doesn't play nicely with this. + // These objects are marshaled into positional arrays on flush. connectionsByProto []map[string]int64 + sessionCounts []map[string]int64 batchSize int // tickCh is used to periodically flush the buffer. @@ -140,6 +141,17 @@ func (b *DBBatcher) Add( st *agentproto.Stats, usage bool, ) { + // Normalize and cap outside the lock. + sessionCounts := normalizedSessionCounts(st) + if len(sessionCounts) > idemetadata.MaxSessionCountEntries { + b.log.Warn(context.Background(), "too many distinct session types, overflow counted under unknown", + slog.F("agent_id", agentID), + slog.F("reported", len(sessionCounts)), + slog.F("max", idemetadata.MaxSessionCountEntries), + ) + } + sessionCounts = capSessionCounts(sessionCounts) + b.mu.Lock() defer b.mu.Unlock() @@ -152,19 +164,14 @@ func (b *DBBatcher) Add( b.buf.TemplateID = append(b.buf.TemplateID, templateID) b.buf.WorkspaceID = append(b.buf.WorkspaceID, workspaceID) - // Store the connections by proto separately as it's a jsonb field. We marshal on flush. - // b.buf.ConnectionsByProto = append(b.buf.ConnectionsByProto, st.ConnectionsByProto) b.connectionsByProto = append(b.connectionsByProto, st.ConnectionsByProto) + b.sessionCounts = append(b.sessionCounts, sessionCounts) b.buf.ConnectionCount = append(b.buf.ConnectionCount, st.ConnectionCount) b.buf.RxPackets = append(b.buf.RxPackets, st.RxPackets) b.buf.RxBytes = append(b.buf.RxBytes, st.RxBytes) b.buf.TxPackets = append(b.buf.TxPackets, st.TxPackets) b.buf.TxBytes = append(b.buf.TxBytes, st.TxBytes) - b.buf.SessionCountVSCode = append(b.buf.SessionCountVSCode, st.SessionCountVscode) - b.buf.SessionCountJetBrains = append(b.buf.SessionCountJetBrains, st.SessionCountJetbrains) - b.buf.SessionCountReconnectingPTY = append(b.buf.SessionCountReconnectingPTY, st.SessionCountReconnectingPty) - b.buf.SessionCountSSH = append(b.buf.SessionCountSSH, st.SessionCountSsh) b.buf.ConnectionMedianLatencyMS = append(b.buf.ConnectionMedianLatencyMS, st.ConnectionMedianLatencyMs) b.buf.Usage = append(b.buf.Usage, usage) @@ -245,6 +252,14 @@ func (b *DBBatcher) flush(ctx context.Context, forced bool, reason string) { b.buf.ConnectionsByProto = payload } + sessionCountsPayload, err := json.Marshal(b.sessionCounts) + if err != nil { + b.log.Error(ctx, "unable to marshal agent session counts, dropping data", slog.Error(err)) + b.buf.SessionCounts = json.RawMessage(`[]`) + } else { + b.buf.SessionCounts = sessionCountsPayload + } + // nolint:gocritic // (#13146) Will be moved soon as part of refactor. err = b.store.InsertWorkspaceAgentStats(ctx, *b.buf) elapsed := time.Since(start) @@ -263,27 +278,25 @@ func (b *DBBatcher) flush(ctx context.Context, forced bool, reason string) { // initBuf resets the buffer. b MUST be locked. func (b *DBBatcher) initBuf(size int) { b.buf = &database.InsertWorkspaceAgentStatsParams{ - ID: make([]uuid.UUID, 0, b.batchSize), - CreatedAt: make([]time.Time, 0, b.batchSize), - UserID: make([]uuid.UUID, 0, b.batchSize), - WorkspaceID: make([]uuid.UUID, 0, b.batchSize), - TemplateID: make([]uuid.UUID, 0, b.batchSize), - AgentID: make([]uuid.UUID, 0, b.batchSize), - ConnectionsByProto: json.RawMessage("[]"), - ConnectionCount: make([]int64, 0, b.batchSize), - RxPackets: make([]int64, 0, b.batchSize), - RxBytes: make([]int64, 0, b.batchSize), - TxPackets: make([]int64, 0, b.batchSize), - TxBytes: make([]int64, 0, b.batchSize), - SessionCountVSCode: make([]int64, 0, b.batchSize), - SessionCountJetBrains: make([]int64, 0, b.batchSize), - SessionCountReconnectingPTY: make([]int64, 0, b.batchSize), - SessionCountSSH: make([]int64, 0, b.batchSize), - ConnectionMedianLatencyMS: make([]float64, 0, b.batchSize), - Usage: make([]bool, 0, b.batchSize), + ID: make([]uuid.UUID, 0, b.batchSize), + CreatedAt: make([]time.Time, 0, b.batchSize), + UserID: make([]uuid.UUID, 0, b.batchSize), + WorkspaceID: make([]uuid.UUID, 0, b.batchSize), + TemplateID: make([]uuid.UUID, 0, b.batchSize), + AgentID: make([]uuid.UUID, 0, b.batchSize), + ConnectionsByProto: json.RawMessage("[]"), + ConnectionCount: make([]int64, 0, b.batchSize), + RxPackets: make([]int64, 0, b.batchSize), + RxBytes: make([]int64, 0, b.batchSize), + TxPackets: make([]int64, 0, b.batchSize), + TxBytes: make([]int64, 0, b.batchSize), + SessionCounts: json.RawMessage("[]"), + ConnectionMedianLatencyMS: make([]float64, 0, b.batchSize), + Usage: make([]bool, 0, b.batchSize), } b.connectionsByProto = make([]map[string]int64, 0, size) + b.sessionCounts = make([]map[string]int64, 0, size) } func (b *DBBatcher) resetBuf() { @@ -299,11 +312,9 @@ func (b *DBBatcher) resetBuf() { b.buf.RxBytes = b.buf.RxBytes[:0] b.buf.TxPackets = b.buf.TxPackets[:0] b.buf.TxBytes = b.buf.TxBytes[:0] - b.buf.SessionCountVSCode = b.buf.SessionCountVSCode[:0] - b.buf.SessionCountJetBrains = b.buf.SessionCountJetBrains[:0] - b.buf.SessionCountReconnectingPTY = b.buf.SessionCountReconnectingPTY[:0] - b.buf.SessionCountSSH = b.buf.SessionCountSSH[:0] + b.buf.SessionCounts = json.RawMessage(`[]`) b.buf.ConnectionMedianLatencyMS = b.buf.ConnectionMedianLatencyMS[:0] b.buf.Usage = b.buf.Usage[:0] b.connectionsByProto = b.connectionsByProto[:0] + b.sessionCounts = b.sessionCounts[:0] } diff --git a/coderd/workspacestats/batcher_internal_test.go b/coderd/workspacestats/batcher_internal_test.go index 48983be561ec3..03ca89fe8d6c9 100644 --- a/coderd/workspacestats/batcher_internal_test.go +++ b/coderd/workspacestats/batcher_internal_test.go @@ -2,9 +2,11 @@ package workspacestats import ( "context" + "encoding/json" "testing" "time" + "github.com/google/uuid" "github.com/stretchr/testify/require" "cdr.dev/slog/v3" @@ -26,7 +28,7 @@ func TestBatchStats(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) t.Cleanup(cancel) log := slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}).Leveled(slog.LevelDebug) - store, ps := dbtestutil.NewDB(t) + store, ps, sqlDB := dbtestutil.NewDBWithSQLDB(t) // Set up some test dependencies. deps1 := setupDeps(t, store, ps) @@ -59,22 +61,54 @@ func TestBatchStats(t *testing.T) { require.NoError(t, err, "should not error getting stats") require.Empty(t, stats, "should have no stats for workspace") - // Given: a single data point is added for workspace + // Given: a stat per workspace, with session counts distinct per agent so + // that a positional misalignment on insert shows up. t2 := t1.Add(time.Second) - t.Log("inserting 1 stat") - b.Add(t2.Add(time.Millisecond), deps1.Agent.ID, deps1.User.ID, deps1.Template.ID, deps1.Workspace.ID, randStats(t), false) + t.Log("inserting 2 stats") + b.Add(t2.Add(time.Millisecond), deps1.Agent.ID, deps1.Template.ID, deps1.User.ID, deps1.Workspace.ID, randStats(t, func(s *agentproto.Stats) { + s.SessionCounts = map[string]int64{"VSCode": 3, "ssh": 1, "idle-ide": 0} + }), false) + b.Add(t2.Add(time.Millisecond), deps2.Agent.ID, deps2.Template.ID, deps2.User.ID, deps2.Workspace.ID, randStats(t, func(s *agentproto.Stats) { + s.SessionCounts = map[string]int64{"jetbrains": 4, "reconnecting-pty": 2} + }), false) // When: it becomes time to report stats // Signal a tick and wait for a flush to complete. tick <- t2 f = <-flushed // Wait for a flush to complete. - require.Equal(t, 1, f, "expected one stat to be flushed") + require.Equal(t, 2, f, "expected two stats to be flushed") t.Log("flush 2 completed") - // Then: it should report a single stat. + // Then: counts reach the right agent, normalized, without the zero entry. stats, err = store.GetWorkspaceAgentStats(ctx, t2) require.NoError(t, err, "should not error getting stats") - require.Len(t, stats, 1, "should have stats for workspace") + require.Len(t, stats, 2, "should have stats for both workspaces") + byAgent := make(map[uuid.UUID]database.GetWorkspaceAgentStatsRow) + for _, stat := range stats { + byAgent[stat.AgentID] = stat + } + require.EqualValues(t, 3, byAgent[deps1.Agent.ID].SessionCountVSCode) + require.EqualValues(t, 1, byAgent[deps1.Agent.ID].SessionCountSSH) + require.EqualValues(t, 0, byAgent[deps1.Agent.ID].SessionCountJetBrains) + require.EqualValues(t, 4, byAgent[deps2.Agent.ID].SessionCountJetBrains) + require.EqualValues(t, 2, byAgent[deps2.Agent.ID].SessionCountReconnectingPTY) + require.EqualValues(t, 0, byAgent[deps2.Agent.ID].SessionCountVSCode) + + // Each row stores its normalized session counts. + storedSessionCounts := func(agentID uuid.UUID) map[string]int64 { + var payload json.RawMessage + require.NoError(t, sqlDB.QueryRowContext(ctx, ` + SELECT session_counts + FROM workspace_agent_stats + WHERE agent_id = $1 AND created_at > $2 + `, agentID, t2).Scan(&payload)) + + var counts map[string]int64 + require.NoError(t, json.Unmarshal(payload, &counts)) + return counts + } + require.Equal(t, map[string]int64{"ssh": 1, "vscode": 3}, storedSessionCounts(deps1.Agent.ID)) + require.Equal(t, map[string]int64{"jetbrains": 4, "reconnecting_pty": 2}, storedSessionCounts(deps2.Agent.ID)) // Given: a lot of data points are added for both workspaces // (equal to batch size) @@ -86,9 +120,9 @@ func TestBatchStats(t *testing.T) { t.Logf("inserting %d stats", defaultBufferSize) for i := 0; i < defaultBufferSize; i++ { if i%2 == 0 { - b.Add(t3.Add(time.Millisecond), deps1.Agent.ID, deps1.User.ID, deps1.Template.ID, deps1.Workspace.ID, randStats(t), false) + b.Add(t3.Add(time.Millisecond), deps1.Agent.ID, deps1.Template.ID, deps1.User.ID, deps1.Workspace.ID, randStats(t), false) } else { - b.Add(t3.Add(time.Millisecond), deps2.Agent.ID, deps2.User.ID, deps2.Template.ID, deps2.Workspace.ID, randStats(t), false) + b.Add(t3.Add(time.Millisecond), deps2.Agent.ID, deps2.Template.ID, deps2.User.ID, deps2.Workspace.ID, randStats(t), false) } } }()