Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion .github/workflows/test-build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,12 @@ jobs:
env:
BILLING_USAGE_TEST_DATABASE_URL: postgresql://postgres:postgres@127.0.0.1:5432/sim_auth_scim
BILLING_USAGE_TEST_REDIS_URL: redis://127.0.0.1:6379
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
run: >-
bunx vitest run
lib/billing/core/usage-log.postgres.test.ts
lib/billing/core/organization-activity.postgres.test.ts
lib/billing/core/usage-analytics-queries.postgres.test.ts
lib/billing/calculations/usage-reservation.test.ts

- name: Verify fork previews ignore execution file history in PostgreSQL
working-directory: apps/sim
Expand Down
68 changes: 44 additions & 24 deletions apps/sim/lib/billing/core/organization-activity-queries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import {
workflowExecutionLogs,
workspace,
} from '@sim/db/schema'
import { and, eq, sql } from 'drizzle-orm'
import { and, eq, type SQL, sql } from 'drizzle-orm'
import {
ACTIVITY_PAGE_SIZE,
type ActivityAggregate,
Expand All @@ -32,11 +32,16 @@ export async function readActivityWorkspace(organizationId: string, workspaceId:
* Only lightweight execution columns are read; transcripts and trace payloads stay private.
* Chat continuations share an execution id and belong to their first retained start.
*/
function activityCte(scope: ActivityScope, dimension?: ActivityDimension) {
function activityGroups(
scope: ActivityScope,
keys: SQL,
grouping: SQL,
dimension?: ActivityDimension
) {
const start = sql`(${scope.start.toISOString()}::timestamptz AT TIME ZONE 'UTC')`
const end = sql`(${scope.end.toISOString()}::timestamptz AT TIME ZONE 'UTC')`
const workflows = sql`
SELECT 'workflow' AS kind, l.workspace_id, l.workflow_id, NULL::text AS member_id,
SELECT l.workspace_id, l.workflow_id,
l.trigger, l.started_at, l.status,
CASE WHEN l.status IN ('completed', 'failed') AND l.total_duration_ms >= 0
THEN l.total_duration_ms END AS duration_ms
Expand All @@ -60,8 +65,7 @@ function activityCte(scope: ActivityScope, dimension?: ActivityDimension) {
WHERE c.organization_id = ${scope.organizationId}`
const chats = sql`
SELECT DISTINCT ON (r.execution_id)
'chat' AS kind, c.workspace_id, NULL::text AS workflow_id, r.user_id AS member_id,
NULL::text AS trigger, r.started_at, NULL::text AS status, NULL::integer AS duration_ms
c.workspace_id, r.user_id AS member_id, r.started_at
FROM ${copilotRuns} r
JOIN (${scopedChats}) c ON c.id = r.chat_id
WHERE r.started_at >= ${start} AND r.started_at < ${end}
Expand All @@ -71,34 +75,45 @@ function activityCte(scope: ActivityScope, dimension?: ActivityDimension) {
)
ORDER BY r.execution_id, r.started_at, r.id
`
const source =
dimension === 'member'
? chats
: dimension === 'workflow' || dimension === 'trigger'
? workflows
: sql`(${workflows}) UNION ALL (${chats})`
return sql`WITH activity AS (${source})`
const workflowGroups = sql`
SELECT ${keys}, count(*) AS "workflowRuns",
count(*) FILTER (WHERE a.status = 'completed') AS completed,
count(*) FILTER (WHERE a.status = 'failed') AS failed,
0::bigint AS "chatRuns", 0::bigint AS "chatMembers",
avg(a.duration_ms) AS "averageDurationMs"
FROM (${workflows}) a ${grouping}
`
const chatGroups = sql`
SELECT ${keys}, 0::bigint AS "workflowRuns", 0::bigint AS completed, 0::bigint AS failed,
count(*) AS "chatRuns", count(DISTINCT a.member_id) AS "chatMembers",
NULL::numeric AS "averageDurationMs"
FROM (${chats}) a ${grouping}
`
if (dimension === 'member') return chatGroups
if (dimension === 'workflow' || dimension === 'trigger') return workflowGroups
return sql`(${workflowGroups}) UNION ALL (${chatGroups})`
}

/** Each group has at most one workflow average and one exact chat-member count. */
const aggregates = sql`
count(*) FILTER (WHERE a.kind = 'workflow') AS "workflowRuns",
count(*) FILTER (WHERE a.kind = 'workflow' AND a.status = 'completed') AS completed,
count(*) FILTER (WHERE a.kind = 'workflow' AND a.status = 'failed') AS failed,
count(*) FILTER (WHERE a.kind = 'chat') AS "chatRuns",
count(DISTINCT a.member_id) AS "chatMembers",
avg(a.duration_ms) AS "averageDurationMs"
sum(a."workflowRuns") AS "workflowRuns",
sum(a.completed) AS completed,
sum(a.failed) AS failed,
sum(a."chatRuns") AS "chatRuns",
sum(a."chatMembers") AS "chatMembers",
max(a."averageDurationMs") AS "averageDurationMs"
`

export async function readActivitySummary(
scope: ActivityScope,
bucket: UsageBucket,
timezone: string
) {
const keys = sql`date_trunc(${bucket}, (a.started_at AT TIME ZONE 'UTC') AT TIME ZONE ${timezone}) AS bucket`
const rows = await dbReplica.execute<ActivityAggregate & { bucket: string | null }>(sql`
${activityCte(scope)}
SELECT to_char(date_trunc(${bucket}, (started_at AT TIME ZONE 'UTC') AT TIME ZONE ${timezone}),
'YYYY-MM-DD') AS bucket, ${aggregates}
FROM activity a GROUP BY GROUPING SETS ((1), ())
WITH activity AS (${activityGroups(scope, keys, sql`GROUP BY GROUPING SETS ((1), ())`)})
SELECT to_char(a.bucket, 'YYYY-MM-DD') AS bucket, ${aggregates}
FROM activity a GROUP BY a.bucket
`)
return {
totals: activityMetrics(rows.find((row) => row.bucket === null)),
Expand Down Expand Up @@ -152,8 +167,13 @@ export async function readActivityBreakdown(
workspaceName: string | null
}
>(sql`
${activityCte(scope, dimension)}, grouped AS (
SELECT ${id} AS id, ${workspaceId} AS "workspaceId", ${aggregates}
WITH activity AS (${activityGroups(
scope,
sql`${id} AS id, ${workspaceId} AS "workspaceId"`,
sql`GROUP BY 1, 2`,
dimension
)}), grouped AS (
SELECT a.id, a."workspaceId", ${aggregates}
FROM activity a
GROUP BY 1, 2
), named AS (
Expand Down
61 changes: 61 additions & 0 deletions apps/sim/lib/billing/core/organization-activity.postgres.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,19 @@ beforeAll(async () => {
('r6', 'c1', 'old', 'm2', '2026-03-09 10:00:00'),
('r7', 'c3', 'foreign', 'm1', '2026-03-09 10:00:00'),
('r8', 'personal', 'personal', 'm1', '2026-03-09 10:00:00');
INSERT INTO workspace VALUES ('edge1', 'First', 'edge'), ('edge2', 'Second', 'edge');
INSERT INTO workflow_execution_logs VALUES
('edge0', 'edge1', 'f1', 'api', '2026-05-01 00:00:00', 'completed', 0),
('edge100', 'edge1', 'f1', 'api', '2026-05-02 00:00:00', 'completed', 100),
('edge300', 'edge2', 'f2', 'manual', '2026-05-02 00:00:00', 'completed', 300),
('negative', 'edge2', 'f2', 'manual', '2026-05-03 00:00:00', 'failed', -1),
('missing', 'edge2', 'f2', 'manual', '2026-05-03 00:00:00', 'completed', NULL);
INSERT INTO copilot_chats VALUES ('edge-chat1', 'edge1', NULL),
('edge-chat2', 'edge2', NULL), ('edge-org-chat', NULL, 'edge');
INSERT INTO copilot_runs VALUES
('edge-r1', 'edge-chat1', 'edge-e1', 'm1', '2026-05-01 00:00:00'),
('edge-r2', 'edge-chat2', 'edge-e2', 'm1', '2026-05-02 00:00:00'),
('edge-r3', 'edge-org-chat', 'edge-e3', 'm1', '2026-05-03 00:00:00');
`)
execute.mockImplementation((query) => database.execute(query))
select.mockImplementation((fields) => database.select(fields))
Expand Down Expand Up @@ -184,4 +197,52 @@ describe.skipIf(!databaseUrl)('organization activity SQL', () => {
averageDurationMs: null,
})
})

it('counts members across the whole period and weights durations by eligible runs', async () => {
const edgeScope = {
organizationId: 'edge',
start: new Date('2026-05-01'),
end: new Date('2026-05-04'),
}
const result = await readActivitySummary(edgeScope, 'day', 'UTC')
expect(result.totals).toMatchObject({
workflowRuns: 5,
completed: 4,
failed: 1,
chatRuns: 3,
chatMembers: 1,
failureRate: 0.2,
})
expect(result.totals.averageDurationMs).toBeCloseTo(400 / 3)
const breakdown = await readActivityBreakdown(edgeScope, 'workspace', 'duration', 0)
expect(breakdown.rows.map((row) => [row.id, row.averageDurationMs, row.chatMembers])).toEqual([
['edge2', 300, 1],
['edge1', 50, 1],
['organization', null, 1],
])
})

it.each(['day', 'week', 'month'] as const)('preserves totals with %s buckets', async (bucket) => {
const result = await readActivitySummary(scope, bucket, 'Pacific/Auckland')
expect(result.totals).toMatchObject({ workflowRuns: 5, chatRuns: 3, chatMembers: 2 })
expect(result.series.reduce((sum, point) => sum + point.workflowRuns, 0)).toBe(5)
expect(result.series.reduce((sum, point) => sum + point.chatRuns, 0)).toBe(3)
})

it('returns workflow-only and chat-only periods without dropping either source', async () => {
const workflowOnly = await readActivitySummary({ ...scope, workspaceId: 'w2' }, 'day', 'UTC')
expect(workflowOnly.totals).toMatchObject({ workflowRuns: 2, chatRuns: 0, chatMembers: 0 })
const chatOnly = await readActivitySummary(
{ ...scope, start: new Date('2026-03-01'), end: new Date('2026-03-02') },
'day',
'UTC'
)
expect(chatOnly.totals).toMatchObject({
workflowRuns: 0,
chatRuns: 1,
chatMembers: 1,
failureRate: null,
averageDurationMs: null,
})
})
})
77 changes: 77 additions & 0 deletions apps/sim/lib/billing/core/usage-analytics-queries.postgres.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
/** @vitest-environment node */
import { generateId } from '@sim/utils/id'
import { eq } from 'drizzle-orm'
import { drizzle } from 'drizzle-orm/postgres-js'
import postgres from 'postgres'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'

const { databaseUrl, select } = vi.hoisted(() => {
const databaseUrl = process.env.BILLING_USAGE_TEST_DATABASE_URL
if (databaseUrl && !['localhost', '127.0.0.1', '[::1]'].includes(new URL(databaseUrl).hostname)) {
throw new Error('Usage integration tests require a disposable local database')
}
return { databaseUrl, select: vi.fn() }
})

vi.unmock('drizzle-orm')
vi.unmock('@sim/db/schema')
vi.mock('@sim/db', () => ({ dbReplica: { select } }))

import { usageLog } from '@sim/db/schema'
import { readUsageTimeSeries } from '@/lib/billing/core/usage-analytics-queries'

const schemaName = `usage_series_${generateId().replaceAll('-', '')}`
const connection = databaseUrl
? postgres(databaseUrl, {
max: 1,
prepare: false,
connection: { search_path: schemaName, timezone: 'Pacific/Auckland' },
onnotice: () => undefined,
})
: undefined

beforeAll(async () => {
if (!connection) return
await connection.unsafe(`CREATE SCHEMA "${schemaName}"`)
await connection.unsafe(`
CREATE TABLE usage_log (billing_entity_id text, created_at timestamp, cost numeric);
INSERT INTO usage_log VALUES
('org', '2026-03-08 08:00:00+00', 0.1),
('org', '2026-03-09 06:59:59+00', 0.2),
('org', '2026-03-09 07:00:00+00', 0.4),
('other', '2026-03-09 07:00:00+00', 999);
`)
const database = drizzle(connection)
select.mockImplementation((fields) => database.select(fields))
})

afterAll(async () => {
if (!connection) return
await connection.unsafe(`DROP SCHEMA "${schemaName}" CASCADE`)
await connection.end()
})

describe.skipIf(!databaseUrl)('usage series SQL', () => {
it('groups the viewer calendar across DST and preserves numeric event counts', async () => {
const rows = await readUsageTimeSeries(
[eq(usageLog.billingEntityId, 'org')],
'day',
'America/Los_Angeles'
)
expect(
rows.toSorted((a, b) => String(a.bucketStart).localeCompare(String(b.bucketStart)))
).toEqual([
{ bucketStart: '2026-03-08T00:00:00', cost: '0.3', events: 2 },
{ bucketStart: '2026-03-09T00:00:00', cost: '0.4', events: 1 },
])
})

it('formats one monthly aggregate and returns no buckets for an empty scope', async () => {
expect(
await readUsageTimeSeries([eq(usageLog.billingEntityId, 'org')], 'month', 'UTC')
).toEqual([{ bucketStart: '2026-03-01T00:00:00', cost: '0.7', events: 3 }])
expect(
await readUsageTimeSeries([eq(usageLog.billingEntityId, 'empty')], 'day', 'UTC')
).toEqual([])
})
})
43 changes: 22 additions & 21 deletions apps/sim/lib/billing/core/usage-analytics-queries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,27 +32,28 @@ export async function readUsageTimeSeries(
executor: DbClient = dbReplica
): Promise<UsageTimeSeriesRow[]> {
assertValidTimezone(timezone)
const bucketStart = sql<string | null>`to_char(
date_trunc(${bucket}, ${usageLog.createdAt} AT TIME ZONE ${timezone}),
'YYYY-MM-DD"T"HH24:MI:SS'
)`

return (
executor
.select({
bucketStart: bucketStart.as('bucket_start'),
cost: sql<string>`COALESCE(SUM(${usageLog.cost}), 0)`,
events: sql<number>`COUNT(*)`.mapWith(Number),
})
.from(usageLog)
.where(and(...scope))
// Group by the output alias, not the expression. Re-rendering the fragment here
// emits a *textually different* one — the select list qualifies the column as
// `created_at`, the group-by as `usage_log.created_at` — and Postgres matches
// group-by expressions syntactically, so it rejects the query outright. It also
// duplicates the bound parameters.
.groupBy(sql`bucket_start`)
)
const buckets = executor
.select({
bucketStart:
sql`date_trunc(${bucket}, (${usageLog.createdAt} AT TIME ZONE 'UTC') AT TIME ZONE ${timezone})`.as(
'bucket_start'
),
cost: sql<string>`COALESCE(SUM(${usageLog.cost}), 0)`.as('cost'),
events: sql<number>`COUNT(*)`.mapWith(Number).as('events'),
})
.from(usageLog)
.where(and(...scope))
.groupBy(sql`bucket_start`)
.as('buckets')

/** Format the aggregated buckets rather than every ledger entry. */
return executor
.select({
bucketStart: sql<string | null>`to_char(${buckets.bucketStart}, 'YYYY-MM-DD"T"HH24:MI:SS')`,
cost: buckets.cost,
events: buckets.events,
})
.from(buckets)
}

export interface UsageTotals {
Expand Down
4 changes: 4 additions & 0 deletions packages/db/db.ts
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,10 @@ export const dbReplica: typeof db = replicaUrl
)
: db

if (!replicaUrl) {
logger.info('Read replica URL is not configured; analytics reads use the primary', { role })
}

const subPoolClients = new Map<SubProcessDbRole, typeof db>()

/** Which env var the process connection came from — named in dbFor fallback logs. */
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
COMMIT;--> statement-breakpoint
SET lock_timeout = 0;--> statement-breakpoint
-- migration-safe: replay replaces only this new index to recover an interrupted concurrent build; existing indexes remain available.
DROP INDEX CONCURRENTLY IF EXISTS "workflow_execution_logs_workspace_activity_idx";--> statement-breakpoint
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
SET lock_timeout = '5s';
Loading
Loading