From 03e17610eae1a29597a2b30eae9bd21b017e6742 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 14:47:09 -0700 Subject: [PATCH 1/2] 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. --- .../mothership/billing/service-delivery.ts | 9 ++-- .../billing/service-store.integration.ts | 54 +++++++++++++++++-- .../lib/mothership/billing/service-store.ts | 49 +++++++++++++---- 3 files changed, 95 insertions(+), 17 deletions(-) diff --git a/apps/sim/lib/mothership/billing/service-delivery.ts b/apps/sim/lib/mothership/billing/service-delivery.ts index d49bcbb01e1..5a44579e53d 100644 --- a/apps/sim/lib/mothership/billing/service-delivery.ts +++ b/apps/sim/lib/mothership/billing/service-delivery.ts @@ -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' @@ -17,8 +17,11 @@ export async function replayServiceUsage(): Promise { 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({ diff --git a/apps/sim/lib/mothership/billing/service-store.integration.ts b/apps/sim/lib/mothership/billing/service-store.integration.ts index 4ba51f49adf..8bbd659f8a5 100644 --- a/apps/sim/lib/mothership/billing/service-store.integration.ts +++ b/apps/sim/lib/mothership/billing/service-store.integration.ts @@ -32,9 +32,9 @@ import { replayServiceUsage } from './service-delivery' import { beginServiceMeter, claimServiceUsage, + closeAbandonedServiceMeters, finishServiceUsage, saveServiceUsage, - serviceMeteringHealth, } from './service-store' afterAll(async () => { @@ -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]) @@ -102,8 +101,53 @@ 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: randomUUID(), + toolCallId: 'abandoned-tool', + workerOrigin: 'http://127.0.0.1:8080', + } + const abandoned = randomUUID() + const failed = randomUUID() + const inFlight = randomUUID() + for (const id of [abandoned, failed, inFlight]) await beginServiceMeter({ ...scope, id }) + await finishServiceUsage(failed, 'provider pricing unavailable') + const receipt = { + id: randomUUID(), + 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([]) + + 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]) }) }) diff --git a/apps/sim/lib/mothership/billing/service-store.ts b/apps/sim/lib/mothership/billing/service-store.ts index 9be46269b66..4f0a64dbc75 100644 --- a/apps/sim/lib/mothership/billing/service-store.ts +++ b/apps/sim/lib/mothership/billing/service-store.ts @@ -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( @@ -59,13 +60,43 @@ export async function beginServiceMeter(input: { .values({ ...input, service: '_tool_execution', costUsd: null }) } -export async function serviceMeteringHealth() { - const [health] = await db - .select({ - pending: sql`count(*) FILTER (WHERE delivered_at IS NULL AND cost_usd IS NOT NULL)::int`, - unknown: sql`count(*) FILTER (WHERE delivered_at IS NULL AND cost_usd IS NULL AND created_at < now() - interval '5 minutes')::int`, - oldest: sql`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 }) + return db + .update(copilotServiceUsage) + .set({ + deliveredAt: new Date(), + lastError: sql`coalesce(${copilotServiceUsage.lastError}, 'Tool execution never finished')`, + }) + .where(inArray(copilotServiceUsage.id, abandoned)) + .returning({ + id: copilotServiceUsage.id, + streamId: copilotServiceUsage.streamId, + toolCallId: copilotServiceUsage.toolCallId, + createdAt: copilotServiceUsage.createdAt, + lastError: copilotServiceUsage.lastError, + }) } From 3072df4c7031989fd11274149e0b2425bf4e8b3e Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 15:34:05 -0700 Subject: [PATCH 2/2] fix(mothership): keep a closed tool meter final against a late tool completion --- .../billing/service-store.integration.ts | 26 ++++++++++++------- .../lib/mothership/billing/service-store.ts | 3 ++- 2 files changed, 19 insertions(+), 10 deletions(-) diff --git a/apps/sim/lib/mothership/billing/service-store.integration.ts b/apps/sim/lib/mothership/billing/service-store.integration.ts index 8bbd659f8a5..8278f8082ca 100644 --- a/apps/sim/lib/mothership/billing/service-store.integration.ts +++ b/apps/sim/lib/mothership/billing/service-store.integration.ts @@ -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' @@ -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', @@ -111,17 +111,17 @@ describe('service receipts in SQL', () => { it('closes each abandoned tool meter once and leaves in-flight meters and receipts open', async () => { const client = state.client! const scope = { - streamId: randomUUID(), + streamId: generateId(), toolCallId: 'abandoned-tool', workerOrigin: 'http://127.0.0.1:8080', } - const abandoned = randomUUID() - const failed = randomUUID() - const inFlight = randomUUID() + 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: randomUUID(), + id: generateId(), streamId: scope.streamId, toolCallId: scope.toolCallId, service: 'exa', @@ -144,6 +144,14 @@ describe('service receipts in SQL', () => { ) 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` diff --git a/apps/sim/lib/mothership/billing/service-store.ts b/apps/sim/lib/mothership/billing/service-store.ts index 4f0a64dbc75..c546c346527 100644 --- a/apps/sim/lib/mothership/billing/service-store.ts +++ b/apps/sim/lib/mothership/billing/service-store.ts @@ -42,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 { 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: {