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
9 changes: 6 additions & 3 deletions apps/sim/lib/mothership/billing/service-delivery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,8 @@ import { getErrorMessage } from '@sim/utils/errors'
import { isHosted } from '@/lib/core/config/env-flags'
import {
claimServiceUsage,
closeAbandonedServiceMeters,
finishServiceUsage,
serviceMeteringHealth,
} from '@/lib/mothership/billing/service-store'
import { ServiceUsageAcknowledgment, ServiceUsageReceipt } from '@/lib/mothership/generated/billing'
import { mothershipRequestHeaders } from '@/lib/mothership/request/headers'
Expand All @@ -17,8 +17,11 @@ export async function replayServiceUsage(): Promise<void> {
if (running) return
running = true
try {
const health = await serviceMeteringHealth()
if (health?.unknown) logger.error('Service usage requires reconciliation', health)
for (const meter of await closeAbandonedServiceMeters())
logger.warn(
'Closed a tool meter that never finished; its provider spend may be unbilled',
meter
)
for (const row of await claimServiceUsage()) {
try {
const receipt = ServiceUsageReceipt.parse({
Expand Down
70 changes: 61 additions & 9 deletions apps/sim/lib/mothership/billing/service-store.integration.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
/** Local SQL verifies receipt durability, concurrent claims and replay independently of tool completion. */

import { randomUUID } from 'node:crypto'
import { readFileSync } from 'node:fs'
import { generateId } from '@sim/utils/id'
import type { Sql } from 'postgres'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'

Expand Down Expand Up @@ -32,9 +32,9 @@ import { replayServiceUsage } from './service-delivery'
import {
beginServiceMeter,
claimServiceUsage,
closeAbandonedServiceMeters,
finishServiceUsage,
saveServiceUsage,
serviceMeteringHealth,
} from './service-store'

afterAll(async () => {
Expand Down Expand Up @@ -68,14 +68,14 @@ describe('service receipts in SQL', () => {
it('claims each known receipt once while preserving incomplete measurement and delivery failures', async () => {
const client = state.client!
const base = {
streamId: randomUUID(),
streamId: generateId(),
toolCallId: 'tool',
workerOrigin: 'http://127.0.0.1:8080',
}
const intentId = randomUUID()
const intentId = generateId()
await beginServiceMeter({ ...base, id: intentId })
const receipt = {
id: randomUUID(),
id: generateId(),
streamId: base.streamId,
toolCallId: base.toolCallId,
service: 'exa',
Expand All @@ -87,8 +87,7 @@ describe('service receipts in SQL', () => {
expect(claims.flat().map((row) => row.id)).toEqual([receipt.id])
await finishServiceUsage(receipt.id, 'connection interrupted')
expect(await claimServiceUsage()).toEqual([])
await client`UPDATE copilot_service_usage SET next_attempt_at=now(), created_at=now()-interval '10 minutes'`
expect((await serviceMeteringHealth())?.unknown).toBe(1)
await client`UPDATE copilot_service_usage SET next_attempt_at=now()`
const fetcher = vi.fn(async (_url: string, options: RequestInit) => {
const body = JSON.parse(String(options.body))
expect(body.receipts).toEqual([receipt])
Expand All @@ -102,8 +101,61 @@ describe('service receipts in SQL', () => {
expect(row.delivered_at).not.toBeNull()
expect(Number(row.cost_usd)).toBe(0.5)
expect(row.worker_origin).toBe(base.workerOrigin)
expect((await serviceMeteringHealth())?.pending).toBe(0)
await finishServiceUsage(intentId)
expect((await serviceMeteringHealth())?.unknown).toBe(0)
const [intent] =
await client`SELECT delivered_at, last_error FROM copilot_service_usage WHERE id=${intentId}`
expect(intent.delivered_at).not.toBeNull()
expect(intent.last_error).toBeNull()
})

it('closes each abandoned tool meter once and leaves in-flight meters and receipts open', async () => {
const client = state.client!
const scope = {
streamId: generateId(),
toolCallId: 'abandoned-tool',
workerOrigin: 'http://127.0.0.1:8080',
}
const abandoned = generateId()
const failed = generateId()
const inFlight = generateId()
for (const id of [abandoned, failed, inFlight]) await beginServiceMeter({ ...scope, id })
await finishServiceUsage(failed, 'provider pricing unavailable')
const receipt = {
id: generateId(),
streamId: scope.streamId,
toolCallId: scope.toolCallId,
service: 'exa',
costUsd: 0.25,
}
await saveServiceUsage(receipt, scope.workerOrigin)
await client`UPDATE copilot_service_usage SET created_at = now() - interval '1 day' WHERE id IN ${client([abandoned, failed, receipt.id])}`
// Past the longest tool watchdog, but a tool can still be cleaning up after it.
await client`UPDATE copilot_service_usage SET created_at = now() - interval '61 minutes' WHERE id = ${inFlight}`

const closed = (
await Promise.all([closeAbandonedServiceMeters(), closeAbandonedServiceMeters()])
).flat()
expect(closed).toHaveLength(2)
expect(new Map(closed.map((meter) => [meter.id, meter.lastError]))).toEqual(
new Map([
[abandoned, expect.any(String)],
[failed, 'provider pricing unavailable'],
])
)
expect(closed.every((meter) => meter.streamId === scope.streamId)).toBe(true)
expect(await closeAbandonedServiceMeters()).toEqual([])
// The watchdog only stops the chat waiting, so the owner can still finish after the close.
await finishServiceUsage(abandoned)
await finishServiceUsage(failed, 'late failure')
const lateRows =
await client`SELECT id, last_error FROM copilot_service_usage WHERE id IN ${client([abandoned, failed])}`
expect(new Map(lateRows.map((row) => [row.id, row.last_error]))).toEqual(
new Map(closed.map((meter) => [meter.id, meter.lastError]))
)

const open =
await client`SELECT id FROM copilot_service_usage WHERE stream_id = ${scope.streamId} AND delivered_at IS NULL`
expect(open.map((row) => row.id).sort()).toEqual([inFlight, receipt.id].sort())
expect((await claimServiceUsage()).map((row) => row.id)).toEqual([receipt.id])
})
})
52 changes: 42 additions & 10 deletions apps/sim/lib/mothership/billing/service-store.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { db } from '@sim/db'
import { copilotServiceUsage } from '@sim/db/schema'
import { and, eq, isNotNull, isNull, lte, sql } from 'drizzle-orm'
import { and, eq, inArray, isNotNull, isNull, lt, lte, sql } from 'drizzle-orm'
import { TOOL_WATCHDOG_LONG_RUNNING_MS } from '@/lib/mothership/constants'
import type { ServiceUsageReceipt } from '@/lib/mothership/generated/billing'

export async function saveServiceUsage(
Expand Down Expand Up @@ -41,11 +42,12 @@ export async function claimServiceUsage(limit = 10) {
})
}

/** A closed row is final, so a tool that outlived its watchdog cannot rewrite its close. */
export async function finishServiceUsage(id: string, error?: string): Promise<void> {
await db
.update(copilotServiceUsage)
.set(error ? { lastError: error } : { deliveredAt: new Date(), lastError: null })
.where(eq(copilotServiceUsage.id, id))
.where(and(eq(copilotServiceUsage.id, id), isNull(copilotServiceUsage.deliveredAt)))
}

export async function beginServiceMeter(input: {
Expand All @@ -59,13 +61,43 @@ export async function beginServiceMeter(input: {
.values({ ...input, service: '_tool_execution', costUsd: null })
}

export async function serviceMeteringHealth() {
const [health] = await db
.select({
pending: sql<number>`count(*) FILTER (WHERE delivered_at IS NULL AND cost_usd IS NOT NULL)::int`,
unknown: sql<number>`count(*) FILTER (WHERE delivered_at IS NULL AND cost_usd IS NULL AND created_at < now() - interval '5 minutes')::int`,
oldest: sql<Date | null>`min(created_at) FILTER (WHERE delivered_at IS NULL)`,
})
/**
* An open meter this old outlived the longest tool watchdog and its cleanup, so the process
* that owned it ended mid-execution and nothing will close it.
*/
const ABANDONED_METER_AGE_MS = 2 * TOOL_WATCHDOG_LONG_RUNNING_MS

/**
* Ends the tool meters that can no longer resolve and returns each one exactly once, so the
* caller reports it once. Closing a meter only ends its audit record: known spend was saved
* as separate receipts, and a meter is never delivered. A pricing failure keeps its error;
* otherwise the row records that the tool never finished.
*/
export async function closeAbandonedServiceMeters(limit = 100) {
const abandoned = db
.select({ id: copilotServiceUsage.id })
.from(copilotServiceUsage)
return health
.where(
and(
isNull(copilotServiceUsage.costUsd),
isNull(copilotServiceUsage.deliveredAt),
lt(copilotServiceUsage.createdAt, new Date(Date.now() - ABANDONED_METER_AGE_MS))
)
)
.limit(limit)
.for('update', { skipLocked: true })
Comment thread
waleedlatif1 marked this conversation as resolved.
return db
.update(copilotServiceUsage)
.set({
deliveredAt: new Date(),
Comment thread
waleedlatif1 marked this conversation as resolved.
lastError: sql`coalesce(${copilotServiceUsage.lastError}, 'Tool execution never finished')`,
Comment thread
waleedlatif1 marked this conversation as resolved.
})
.where(inArray(copilotServiceUsage.id, abandoned))
.returning({
id: copilotServiceUsage.id,
streamId: copilotServiceUsage.streamId,
toolCallId: copilotServiceUsage.toolCallId,
createdAt: copilotServiceUsage.createdAt,
lastError: copilotServiceUsage.lastError,
})
}
Loading