Skip to content

Commit f2a89aa

Browse files
authored
fix(insights): reduce activity aggregation overhead (#7908)
* fix(insights): reduce activity aggregation overhead * fix(db): declare activity index as concurrent
1 parent f7f678b commit f2a89aa

10 files changed

Lines changed: 27525 additions & 46 deletions

‎.github/workflows/test-build.yml‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -119,7 +119,12 @@ jobs:
119119
env:
120120
BILLING_USAGE_TEST_DATABASE_URL: postgresql://postgres:postgres@127.0.0.1:5432/sim_auth_scim
121121
BILLING_USAGE_TEST_REDIS_URL: redis://127.0.0.1:6379
122-
run: bunx vitest run lib/billing/core/usage-log.postgres.test.ts lib/billing/core/organization-activity.postgres.test.ts lib/billing/calculations/usage-reservation.test.ts
122+
run: >-
123+
bunx vitest run
124+
lib/billing/core/usage-log.postgres.test.ts
125+
lib/billing/core/organization-activity.postgres.test.ts
126+
lib/billing/core/usage-analytics-queries.postgres.test.ts
127+
lib/billing/calculations/usage-reservation.test.ts
123128
124129
- name: Verify fork previews ignore execution file history in PostgreSQL
125130
working-directory: apps/sim

‎apps/sim/lib/billing/core/organization-activity-queries.ts‎

Lines changed: 44 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ import {
77
workflowExecutionLogs,
88
workspace,
99
} from '@sim/db/schema'
10-
import { and, eq, sql } from 'drizzle-orm'
10+
import { and, eq, type SQL, sql } from 'drizzle-orm'
1111
import {
1212
ACTIVITY_PAGE_SIZE,
1313
type ActivityAggregate,
@@ -32,11 +32,16 @@ export async function readActivityWorkspace(organizationId: string, workspaceId:
3232
* Only lightweight execution columns are read; transcripts and trace payloads stay private.
3333
* Chat continuations share an execution id and belong to their first retained start.
3434
*/
35-
function activityCte(scope: ActivityScope, dimension?: ActivityDimension) {
35+
function activityGroups(
36+
scope: ActivityScope,
37+
keys: SQL,
38+
grouping: SQL,
39+
dimension?: ActivityDimension
40+
) {
3641
const start = sql`(${scope.start.toISOString()}::timestamptz AT TIME ZONE 'UTC')`
3742
const end = sql`(${scope.end.toISOString()}::timestamptz AT TIME ZONE 'UTC')`
3843
const workflows = sql`
39-
SELECT 'workflow' AS kind, l.workspace_id, l.workflow_id, NULL::text AS member_id,
44+
SELECT l.workspace_id, l.workflow_id,
4045
l.trigger, l.started_at, l.status,
4146
CASE WHEN l.status IN ('completed', 'failed') AND l.total_duration_ms >= 0
4247
THEN l.total_duration_ms END AS duration_ms
@@ -60,8 +65,7 @@ function activityCte(scope: ActivityScope, dimension?: ActivityDimension) {
6065
WHERE c.organization_id = ${scope.organizationId}`
6166
const chats = sql`
6267
SELECT DISTINCT ON (r.execution_id)
63-
'chat' AS kind, c.workspace_id, NULL::text AS workflow_id, r.user_id AS member_id,
64-
NULL::text AS trigger, r.started_at, NULL::text AS status, NULL::integer AS duration_ms
68+
c.workspace_id, r.user_id AS member_id, r.started_at
6569
FROM ${copilotRuns} r
6670
JOIN (${scopedChats}) c ON c.id = r.chat_id
6771
WHERE r.started_at >= ${start} AND r.started_at < ${end}
@@ -71,34 +75,45 @@ function activityCte(scope: ActivityScope, dimension?: ActivityDimension) {
7175
)
7276
ORDER BY r.execution_id, r.started_at, r.id
7377
`
74-
const source =
75-
dimension === 'member'
76-
? chats
77-
: dimension === 'workflow' || dimension === 'trigger'
78-
? workflows
79-
: sql`(${workflows}) UNION ALL (${chats})`
80-
return sql`WITH activity AS (${source})`
78+
const workflowGroups = sql`
79+
SELECT ${keys}, count(*) AS "workflowRuns",
80+
count(*) FILTER (WHERE a.status = 'completed') AS completed,
81+
count(*) FILTER (WHERE a.status = 'failed') AS failed,
82+
0::bigint AS "chatRuns", 0::bigint AS "chatMembers",
83+
avg(a.duration_ms) AS "averageDurationMs"
84+
FROM (${workflows}) a ${grouping}
85+
`
86+
const chatGroups = sql`
87+
SELECT ${keys}, 0::bigint AS "workflowRuns", 0::bigint AS completed, 0::bigint AS failed,
88+
count(*) AS "chatRuns", count(DISTINCT a.member_id) AS "chatMembers",
89+
NULL::numeric AS "averageDurationMs"
90+
FROM (${chats}) a ${grouping}
91+
`
92+
if (dimension === 'member') return chatGroups
93+
if (dimension === 'workflow' || dimension === 'trigger') return workflowGroups
94+
return sql`(${workflowGroups}) UNION ALL (${chatGroups})`
8195
}
8296

97+
/** Each group has at most one workflow average and one exact chat-member count. */
8398
const aggregates = sql`
84-
count(*) FILTER (WHERE a.kind = 'workflow') AS "workflowRuns",
85-
count(*) FILTER (WHERE a.kind = 'workflow' AND a.status = 'completed') AS completed,
86-
count(*) FILTER (WHERE a.kind = 'workflow' AND a.status = 'failed') AS failed,
87-
count(*) FILTER (WHERE a.kind = 'chat') AS "chatRuns",
88-
count(DISTINCT a.member_id) AS "chatMembers",
89-
avg(a.duration_ms) AS "averageDurationMs"
99+
sum(a."workflowRuns") AS "workflowRuns",
100+
sum(a.completed) AS completed,
101+
sum(a.failed) AS failed,
102+
sum(a."chatRuns") AS "chatRuns",
103+
sum(a."chatMembers") AS "chatMembers",
104+
max(a."averageDurationMs") AS "averageDurationMs"
90105
`
91106

92107
export async function readActivitySummary(
93108
scope: ActivityScope,
94109
bucket: UsageBucket,
95110
timezone: string
96111
) {
112+
const keys = sql`date_trunc(${bucket}, (a.started_at AT TIME ZONE 'UTC') AT TIME ZONE ${timezone}) AS bucket`
97113
const rows = await dbReplica.execute<ActivityAggregate & { bucket: string | null }>(sql`
98-
${activityCte(scope)}
99-
SELECT to_char(date_trunc(${bucket}, (started_at AT TIME ZONE 'UTC') AT TIME ZONE ${timezone}),
100-
'YYYY-MM-DD') AS bucket, ${aggregates}
101-
FROM activity a GROUP BY GROUPING SETS ((1), ())
114+
WITH activity AS (${activityGroups(scope, keys, sql`GROUP BY GROUPING SETS ((1), ())`)})
115+
SELECT to_char(a.bucket, 'YYYY-MM-DD') AS bucket, ${aggregates}
116+
FROM activity a GROUP BY a.bucket
102117
`)
103118
return {
104119
totals: activityMetrics(rows.find((row) => row.bucket === null)),
@@ -152,8 +167,13 @@ export async function readActivityBreakdown(
152167
workspaceName: string | null
153168
}
154169
>(sql`
155-
${activityCte(scope, dimension)}, grouped AS (
156-
SELECT ${id} AS id, ${workspaceId} AS "workspaceId", ${aggregates}
170+
WITH activity AS (${activityGroups(
171+
scope,
172+
sql`${id} AS id, ${workspaceId} AS "workspaceId"`,
173+
sql`GROUP BY 1, 2`,
174+
dimension
175+
)}), grouped AS (
176+
SELECT a.id, a."workspaceId", ${aggregates}
157177
FROM activity a
158178
GROUP BY 1, 2
159179
), named AS (

‎apps/sim/lib/billing/core/organization-activity.postgres.test.ts‎

Lines changed: 61 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,19 @@ beforeAll(async () => {
7777
('r6', 'c1', 'old', 'm2', '2026-03-09 10:00:00'),
7878
('r7', 'c3', 'foreign', 'm1', '2026-03-09 10:00:00'),
7979
('r8', 'personal', 'personal', 'm1', '2026-03-09 10:00:00');
80+
INSERT INTO workspace VALUES ('edge1', 'First', 'edge'), ('edge2', 'Second', 'edge');
81+
INSERT INTO workflow_execution_logs VALUES
82+
('edge0', 'edge1', 'f1', 'api', '2026-05-01 00:00:00', 'completed', 0),
83+
('edge100', 'edge1', 'f1', 'api', '2026-05-02 00:00:00', 'completed', 100),
84+
('edge300', 'edge2', 'f2', 'manual', '2026-05-02 00:00:00', 'completed', 300),
85+
('negative', 'edge2', 'f2', 'manual', '2026-05-03 00:00:00', 'failed', -1),
86+
('missing', 'edge2', 'f2', 'manual', '2026-05-03 00:00:00', 'completed', NULL);
87+
INSERT INTO copilot_chats VALUES ('edge-chat1', 'edge1', NULL),
88+
('edge-chat2', 'edge2', NULL), ('edge-org-chat', NULL, 'edge');
89+
INSERT INTO copilot_runs VALUES
90+
('edge-r1', 'edge-chat1', 'edge-e1', 'm1', '2026-05-01 00:00:00'),
91+
('edge-r2', 'edge-chat2', 'edge-e2', 'm1', '2026-05-02 00:00:00'),
92+
('edge-r3', 'edge-org-chat', 'edge-e3', 'm1', '2026-05-03 00:00:00');
8093
`)
8194
execute.mockImplementation((query) => database.execute(query))
8295
select.mockImplementation((fields) => database.select(fields))
@@ -184,4 +197,52 @@ describe.skipIf(!databaseUrl)('organization activity SQL', () => {
184197
averageDurationMs: null,
185198
})
186199
})
200+
201+
it('counts members across the whole period and weights durations by eligible runs', async () => {
202+
const edgeScope = {
203+
organizationId: 'edge',
204+
start: new Date('2026-05-01'),
205+
end: new Date('2026-05-04'),
206+
}
207+
const result = await readActivitySummary(edgeScope, 'day', 'UTC')
208+
expect(result.totals).toMatchObject({
209+
workflowRuns: 5,
210+
completed: 4,
211+
failed: 1,
212+
chatRuns: 3,
213+
chatMembers: 1,
214+
failureRate: 0.2,
215+
})
216+
expect(result.totals.averageDurationMs).toBeCloseTo(400 / 3)
217+
const breakdown = await readActivityBreakdown(edgeScope, 'workspace', 'duration', 0)
218+
expect(breakdown.rows.map((row) => [row.id, row.averageDurationMs, row.chatMembers])).toEqual([
219+
['edge2', 300, 1],
220+
['edge1', 50, 1],
221+
['organization', null, 1],
222+
])
223+
})
224+
225+
it.each(['day', 'week', 'month'] as const)('preserves totals with %s buckets', async (bucket) => {
226+
const result = await readActivitySummary(scope, bucket, 'Pacific/Auckland')
227+
expect(result.totals).toMatchObject({ workflowRuns: 5, chatRuns: 3, chatMembers: 2 })
228+
expect(result.series.reduce((sum, point) => sum + point.workflowRuns, 0)).toBe(5)
229+
expect(result.series.reduce((sum, point) => sum + point.chatRuns, 0)).toBe(3)
230+
})
231+
232+
it('returns workflow-only and chat-only periods without dropping either source', async () => {
233+
const workflowOnly = await readActivitySummary({ ...scope, workspaceId: 'w2' }, 'day', 'UTC')
234+
expect(workflowOnly.totals).toMatchObject({ workflowRuns: 2, chatRuns: 0, chatMembers: 0 })
235+
const chatOnly = await readActivitySummary(
236+
{ ...scope, start: new Date('2026-03-01'), end: new Date('2026-03-02') },
237+
'day',
238+
'UTC'
239+
)
240+
expect(chatOnly.totals).toMatchObject({
241+
workflowRuns: 0,
242+
chatRuns: 1,
243+
chatMembers: 1,
244+
failureRate: null,
245+
averageDurationMs: null,
246+
})
247+
})
187248
})
Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
/** @vitest-environment node */
2+
import { generateId } from '@sim/utils/id'
3+
import { eq } from 'drizzle-orm'
4+
import { drizzle } from 'drizzle-orm/postgres-js'
5+
import postgres from 'postgres'
6+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
7+
8+
const { databaseUrl, select } = vi.hoisted(() => {
9+
const databaseUrl = process.env.BILLING_USAGE_TEST_DATABASE_URL
10+
if (databaseUrl && !['localhost', '127.0.0.1', '[::1]'].includes(new URL(databaseUrl).hostname)) {
11+
throw new Error('Usage integration tests require a disposable local database')
12+
}
13+
return { databaseUrl, select: vi.fn() }
14+
})
15+
16+
vi.unmock('drizzle-orm')
17+
vi.unmock('@sim/db/schema')
18+
vi.mock('@sim/db', () => ({ dbReplica: { select } }))
19+
20+
import { usageLog } from '@sim/db/schema'
21+
import { readUsageTimeSeries } from '@/lib/billing/core/usage-analytics-queries'
22+
23+
const schemaName = `usage_series_${generateId().replaceAll('-', '')}`
24+
const connection = databaseUrl
25+
? postgres(databaseUrl, {
26+
max: 1,
27+
prepare: false,
28+
connection: { search_path: schemaName, timezone: 'Pacific/Auckland' },
29+
onnotice: () => undefined,
30+
})
31+
: undefined
32+
33+
beforeAll(async () => {
34+
if (!connection) return
35+
await connection.unsafe(`CREATE SCHEMA "${schemaName}"`)
36+
await connection.unsafe(`
37+
CREATE TABLE usage_log (billing_entity_id text, created_at timestamp, cost numeric);
38+
INSERT INTO usage_log VALUES
39+
('org', '2026-03-08 08:00:00+00', 0.1),
40+
('org', '2026-03-09 06:59:59+00', 0.2),
41+
('org', '2026-03-09 07:00:00+00', 0.4),
42+
('other', '2026-03-09 07:00:00+00', 999);
43+
`)
44+
const database = drizzle(connection)
45+
select.mockImplementation((fields) => database.select(fields))
46+
})
47+
48+
afterAll(async () => {
49+
if (!connection) return
50+
await connection.unsafe(`DROP SCHEMA "${schemaName}" CASCADE`)
51+
await connection.end()
52+
})
53+
54+
describe.skipIf(!databaseUrl)('usage series SQL', () => {
55+
it('groups the viewer calendar across DST and preserves numeric event counts', async () => {
56+
const rows = await readUsageTimeSeries(
57+
[eq(usageLog.billingEntityId, 'org')],
58+
'day',
59+
'America/Los_Angeles'
60+
)
61+
expect(
62+
rows.toSorted((a, b) => String(a.bucketStart).localeCompare(String(b.bucketStart)))
63+
).toEqual([
64+
{ bucketStart: '2026-03-08T00:00:00', cost: '0.3', events: 2 },
65+
{ bucketStart: '2026-03-09T00:00:00', cost: '0.4', events: 1 },
66+
])
67+
})
68+
69+
it('formats one monthly aggregate and returns no buckets for an empty scope', async () => {
70+
expect(
71+
await readUsageTimeSeries([eq(usageLog.billingEntityId, 'org')], 'month', 'UTC')
72+
).toEqual([{ bucketStart: '2026-03-01T00:00:00', cost: '0.7', events: 3 }])
73+
expect(
74+
await readUsageTimeSeries([eq(usageLog.billingEntityId, 'empty')], 'day', 'UTC')
75+
).toEqual([])
76+
})
77+
})

‎apps/sim/lib/billing/core/usage-analytics-queries.ts‎

Lines changed: 22 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -32,27 +32,28 @@ export async function readUsageTimeSeries(
3232
executor: DbClient = dbReplica
3333
): Promise<UsageTimeSeriesRow[]> {
3434
assertValidTimezone(timezone)
35-
const bucketStart = sql<string | null>`to_char(
36-
date_trunc(${bucket}, ${usageLog.createdAt} AT TIME ZONE ${timezone}),
37-
'YYYY-MM-DD"T"HH24:MI:SS'
38-
)`
39-
40-
return (
41-
executor
42-
.select({
43-
bucketStart: bucketStart.as('bucket_start'),
44-
cost: sql<string>`COALESCE(SUM(${usageLog.cost}), 0)`,
45-
events: sql<number>`COUNT(*)`.mapWith(Number),
46-
})
47-
.from(usageLog)
48-
.where(and(...scope))
49-
// Group by the output alias, not the expression. Re-rendering the fragment here
50-
// emits a *textually different* one — the select list qualifies the column as
51-
// `created_at`, the group-by as `usage_log.created_at` — and Postgres matches
52-
// group-by expressions syntactically, so it rejects the query outright. It also
53-
// duplicates the bound parameters.
54-
.groupBy(sql`bucket_start`)
55-
)
35+
const buckets = executor
36+
.select({
37+
bucketStart:
38+
sql`date_trunc(${bucket}, (${usageLog.createdAt} AT TIME ZONE 'UTC') AT TIME ZONE ${timezone})`.as(
39+
'bucket_start'
40+
),
41+
cost: sql<string>`COALESCE(SUM(${usageLog.cost}), 0)`.as('cost'),
42+
events: sql<number>`COUNT(*)`.mapWith(Number).as('events'),
43+
})
44+
.from(usageLog)
45+
.where(and(...scope))
46+
.groupBy(sql`bucket_start`)
47+
.as('buckets')
48+
49+
/** Format the aggregated buckets rather than every ledger entry. */
50+
return executor
51+
.select({
52+
bucketStart: sql<string | null>`to_char(${buckets.bucketStart}, 'YYYY-MM-DD"T"HH24:MI:SS')`,
53+
cost: buckets.cost,
54+
events: buckets.events,
55+
})
56+
.from(buckets)
5657
}
5758

5859
export interface UsageTotals {

‎packages/db/db.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -113,6 +113,10 @@ export const dbReplica: typeof db = replicaUrl
113113
)
114114
: db
115115

116+
if (!replicaUrl) {
117+
logger.info('Read replica URL is not configured; analytics reads use the primary', { role })
118+
}
119+
116120
const subPoolClients = new Map<SubProcessDbRole, typeof db>()
117121

118122
/** Which env var the process connection came from — named in dbFor fallback logs. */
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
COMMIT;--> statement-breakpoint
2+
SET lock_timeout = 0;--> statement-breakpoint
3+
-- migration-safe: replay replaces only this new index to recover an interrupted concurrent build; existing indexes remain available.
4+
DROP INDEX CONCURRENTLY IF EXISTS "workflow_execution_logs_workspace_activity_idx";--> statement-breakpoint
5+
CREATE INDEX CONCURRENTLY IF NOT EXISTS "workflow_execution_logs_workspace_activity_idx" ON "workflow_execution_logs" USING btree ("workspace_id","started_at","status","total_duration_ms","workflow_id","trigger");--> statement-breakpoint
6+
SET lock_timeout = '5s';

0 commit comments

Comments
 (0)