Skip to content

Commit fb33da1

Browse files
authored
fix(mothership): close abandoned tool meters once instead of alarming every tick (#8479)
* fix(mothership): close abandoned tool meters once instead of alarming every tick A tool meter row (cost unknown) stays open when the process that owned the tool ends mid-execution, and nothing ever closed it. The replay tick counted those rows and logged "Service usage requires reconciliation" at ERROR on every tick in every process, forever. Its 5-minute threshold also flagged tools that were still legitimately running. The replay tick now closes meters older than twice the longest tool watchdog, keeping a pricing failure's error or recording that the tool never finished, and logs each closed meter once with its stream, tool call and reason. Known spend is unaffected: it is saved and delivered as separate receipts. * fix(mothership): keep a closed tool meter final against a late tool completion
1 parent f89b168 commit fb33da1

3 files changed

Lines changed: 109 additions & 22 deletions

File tree

‎apps/sim/lib/mothership/billing/service-delivery.ts‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,8 @@ import { getErrorMessage } from '@sim/utils/errors'
33
import { isHosted } from '@/lib/core/config/env-flags'
44
import {
55
claimServiceUsage,
6+
closeAbandonedServiceMeters,
67
finishServiceUsage,
7-
serviceMeteringHealth,
88
} from '@/lib/mothership/billing/service-store'
99
import { ServiceUsageAcknowledgment, ServiceUsageReceipt } from '@/lib/mothership/generated/billing'
1010
import { mothershipRequestHeaders } from '@/lib/mothership/request/headers'
@@ -17,8 +17,11 @@ export async function replayServiceUsage(): Promise<void> {
1717
if (running) return
1818
running = true
1919
try {
20-
const health = await serviceMeteringHealth()
21-
if (health?.unknown) logger.error('Service usage requires reconciliation', health)
20+
for (const meter of await closeAbandonedServiceMeters())
21+
logger.warn(
22+
'Closed a tool meter that never finished; its provider spend may be unbilled',
23+
meter
24+
)
2225
for (const row of await claimServiceUsage()) {
2326
try {
2427
const receipt = ServiceUsageReceipt.parse({

‎apps/sim/lib/mothership/billing/service-store.integration.ts‎

Lines changed: 61 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
/** Local SQL verifies receipt durability, concurrent claims and replay independently of tool completion. */
22

3-
import { randomUUID } from 'node:crypto'
43
import { readFileSync } from 'node:fs'
4+
import { generateId } from '@sim/utils/id'
55
import type { Sql } from 'postgres'
66
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
77

@@ -32,9 +32,9 @@ import { replayServiceUsage } from './service-delivery'
3232
import {
3333
beginServiceMeter,
3434
claimServiceUsage,
35+
closeAbandonedServiceMeters,
3536
finishServiceUsage,
3637
saveServiceUsage,
37-
serviceMeteringHealth,
3838
} from './service-store'
3939

4040
afterAll(async () => {
@@ -68,14 +68,14 @@ describe('service receipts in SQL', () => {
6868
it('claims each known receipt once while preserving incomplete measurement and delivery failures', async () => {
6969
const client = state.client!
7070
const base = {
71-
streamId: randomUUID(),
71+
streamId: generateId(),
7272
toolCallId: 'tool',
7373
workerOrigin: 'http://127.0.0.1:8080',
7474
}
75-
const intentId = randomUUID()
75+
const intentId = generateId()
7676
await beginServiceMeter({ ...base, id: intentId })
7777
const receipt = {
78-
id: randomUUID(),
78+
id: generateId(),
7979
streamId: base.streamId,
8080
toolCallId: base.toolCallId,
8181
service: 'exa',
@@ -87,8 +87,7 @@ describe('service receipts in SQL', () => {
8787
expect(claims.flat().map((row) => row.id)).toEqual([receipt.id])
8888
await finishServiceUsage(receipt.id, 'connection interrupted')
8989
expect(await claimServiceUsage()).toEqual([])
90-
await client`UPDATE copilot_service_usage SET next_attempt_at=now(), created_at=now()-interval '10 minutes'`
91-
expect((await serviceMeteringHealth())?.unknown).toBe(1)
90+
await client`UPDATE copilot_service_usage SET next_attempt_at=now()`
9291
const fetcher = vi.fn(async (_url: string, options: RequestInit) => {
9392
const body = JSON.parse(String(options.body))
9493
expect(body.receipts).toEqual([receipt])
@@ -102,8 +101,61 @@ describe('service receipts in SQL', () => {
102101
expect(row.delivered_at).not.toBeNull()
103102
expect(Number(row.cost_usd)).toBe(0.5)
104103
expect(row.worker_origin).toBe(base.workerOrigin)
105-
expect((await serviceMeteringHealth())?.pending).toBe(0)
106104
await finishServiceUsage(intentId)
107-
expect((await serviceMeteringHealth())?.unknown).toBe(0)
105+
const [intent] =
106+
await client`SELECT delivered_at, last_error FROM copilot_service_usage WHERE id=${intentId}`
107+
expect(intent.delivered_at).not.toBeNull()
108+
expect(intent.last_error).toBeNull()
109+
})
110+
111+
it('closes each abandoned tool meter once and leaves in-flight meters and receipts open', async () => {
112+
const client = state.client!
113+
const scope = {
114+
streamId: generateId(),
115+
toolCallId: 'abandoned-tool',
116+
workerOrigin: 'http://127.0.0.1:8080',
117+
}
118+
const abandoned = generateId()
119+
const failed = generateId()
120+
const inFlight = generateId()
121+
for (const id of [abandoned, failed, inFlight]) await beginServiceMeter({ ...scope, id })
122+
await finishServiceUsage(failed, 'provider pricing unavailable')
123+
const receipt = {
124+
id: generateId(),
125+
streamId: scope.streamId,
126+
toolCallId: scope.toolCallId,
127+
service: 'exa',
128+
costUsd: 0.25,
129+
}
130+
await saveServiceUsage(receipt, scope.workerOrigin)
131+
await client`UPDATE copilot_service_usage SET created_at = now() - interval '1 day' WHERE id IN ${client([abandoned, failed, receipt.id])}`
132+
// Past the longest tool watchdog, but a tool can still be cleaning up after it.
133+
await client`UPDATE copilot_service_usage SET created_at = now() - interval '61 minutes' WHERE id = ${inFlight}`
134+
135+
const closed = (
136+
await Promise.all([closeAbandonedServiceMeters(), closeAbandonedServiceMeters()])
137+
).flat()
138+
expect(closed).toHaveLength(2)
139+
expect(new Map(closed.map((meter) => [meter.id, meter.lastError]))).toEqual(
140+
new Map([
141+
[abandoned, expect.any(String)],
142+
[failed, 'provider pricing unavailable'],
143+
])
144+
)
145+
expect(closed.every((meter) => meter.streamId === scope.streamId)).toBe(true)
146+
expect(await closeAbandonedServiceMeters()).toEqual([])
147+
// The watchdog only stops the chat waiting, so the owner can still finish after the close.
148+
await finishServiceUsage(abandoned)
149+
await finishServiceUsage(failed, 'late failure')
150+
const lateRows =
151+
await client`SELECT id, last_error FROM copilot_service_usage WHERE id IN ${client([abandoned, failed])}`
152+
expect(new Map(lateRows.map((row) => [row.id, row.last_error]))).toEqual(
153+
new Map(closed.map((meter) => [meter.id, meter.lastError]))
154+
)
155+
156+
const open =
157+
await client`SELECT id FROM copilot_service_usage WHERE stream_id = ${scope.streamId} AND delivered_at IS NULL`
158+
expect(open.map((row) => row.id).sort()).toEqual([inFlight, receipt.id].sort())
159+
expect((await claimServiceUsage()).map((row) => row.id)).toEqual([receipt.id])
108160
})
109161
})

‎apps/sim/lib/mothership/billing/service-store.ts‎

Lines changed: 42 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { db } from '@sim/db'
22
import { copilotServiceUsage } from '@sim/db/schema'
3-
import { and, eq, isNotNull, isNull, lte, sql } from 'drizzle-orm'
3+
import { and, eq, inArray, isNotNull, isNull, lt, lte, sql } from 'drizzle-orm'
4+
import { TOOL_WATCHDOG_LONG_RUNNING_MS } from '@/lib/mothership/constants'
45
import type { ServiceUsageReceipt } from '@/lib/mothership/generated/billing'
56

67
export async function saveServiceUsage(
@@ -41,11 +42,12 @@ export async function claimServiceUsage(limit = 10) {
4142
})
4243
}
4344

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

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

62-
export async function serviceMeteringHealth() {
63-
const [health] = await db
64-
.select({
65-
pending: sql<number>`count(*) FILTER (WHERE delivered_at IS NULL AND cost_usd IS NOT NULL)::int`,
66-
unknown: sql<number>`count(*) FILTER (WHERE delivered_at IS NULL AND cost_usd IS NULL AND created_at < now() - interval '5 minutes')::int`,
67-
oldest: sql<Date | null>`min(created_at) FILTER (WHERE delivered_at IS NULL)`,
68-
})
64+
/**
65+
* An open meter this old outlived the longest tool watchdog and its cleanup, so the process
66+
* that owned it ended mid-execution and nothing will close it.
67+
*/
68+
const ABANDONED_METER_AGE_MS = 2 * TOOL_WATCHDOG_LONG_RUNNING_MS
69+
70+
/**
71+
* Ends the tool meters that can no longer resolve and returns each one exactly once, so the
72+
* caller reports it once. Closing a meter only ends its audit record: known spend was saved
73+
* as separate receipts, and a meter is never delivered. A pricing failure keeps its error;
74+
* otherwise the row records that the tool never finished.
75+
*/
76+
export async function closeAbandonedServiceMeters(limit = 100) {
77+
const abandoned = db
78+
.select({ id: copilotServiceUsage.id })
6979
.from(copilotServiceUsage)
70-
return health
80+
.where(
81+
and(
82+
isNull(copilotServiceUsage.costUsd),
83+
isNull(copilotServiceUsage.deliveredAt),
84+
lt(copilotServiceUsage.createdAt, new Date(Date.now() - ABANDONED_METER_AGE_MS))
85+
)
86+
)
87+
.limit(limit)
88+
.for('update', { skipLocked: true })
89+
return db
90+
.update(copilotServiceUsage)
91+
.set({
92+
deliveredAt: new Date(),
93+
lastError: sql`coalesce(${copilotServiceUsage.lastError}, 'Tool execution never finished')`,
94+
})
95+
.where(inArray(copilotServiceUsage.id, abandoned))
96+
.returning({
97+
id: copilotServiceUsage.id,
98+
streamId: copilotServiceUsage.streamId,
99+
toolCallId: copilotServiceUsage.toolCallId,
100+
createdAt: copilotServiceUsage.createdAt,
101+
lastError: copilotServiceUsage.lastError,
102+
})
71103
}

0 commit comments

Comments
 (0)