Skip to content

Commit 77273e3

Browse files
committed
fix(mothership): preserve hosted service billing through cutover
Record hosted integration, sandbox and media spend through a trusted durable service outbox independent of tool output, cancellation and worker delivery. Preserve service billing under model BYOK and restore title admission. Use workspace BYOK settings and fresh credentials, retain execute transport, and restore the locked actor-authority check for billing-sensitive membership removal. Add migrations and local HTTP/PostgreSQL billing acceptance coverage. Validation: 1,453 billing tests passed with 30 environment-gated skips, app and dev infrastructure types passed, and local SQL charge proofs passed.
1 parent fef222f commit 77273e3

45 files changed

Lines changed: 58810 additions & 358 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
Lines changed: 206 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,206 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { type ExecFileException, execFile } from 'node:child_process'
5+
import { createServer } from 'node:http'
6+
import { promisify } from 'node:util'
7+
import { resetEnvFlagsMock, resetEnvMock, setEnv, setEnvFlags } from '@sim/testing'
8+
import { NextRequest } from 'next/server'
9+
import type { Sql } from 'postgres'
10+
import { afterAll, describe, expect, it, vi } from 'vitest'
11+
12+
const state = vi.hoisted(() => ({
13+
databaseUrl: process.env.BILLING_REPLAY_IT_DATABASE_URL,
14+
workerDirectory: process.env.BILLING_REPLAY_WORKER_DIR,
15+
client: null as Sql | null,
16+
schema: `billing_worker_${process.pid}`,
17+
temporaryFailures: 1,
18+
}))
19+
20+
vi.unmock('drizzle-orm')
21+
vi.unmock('@sim/db/schema')
22+
vi.mock('@sim/db', async () => {
23+
const { drizzle } = await import('drizzle-orm/postgres-js')
24+
const { default: postgres } = await import('postgres')
25+
const client = postgres(state.databaseUrl ?? 'postgres://127.0.0.1:1/unused', {
26+
max: 2,
27+
connection: { search_path: state.schema },
28+
onnotice: () => {},
29+
})
30+
state.client = client
31+
const db = drizzle(client)
32+
return { db, dbReplica: db }
33+
})
34+
35+
/** Subscription and payment-provider fixtures; callback, ledger and replay code are real. */
36+
vi.mock('@/lib/billing/core/subscription', () => {
37+
const subscription = async (userId: string) => {
38+
if (userId === 'billing-replay-transient' && state.temporaryFailures-- > 0) {
39+
throw new Error('Temporary subscription lookup failure')
40+
}
41+
return {
42+
id: 'billing-replay-subscription',
43+
referenceId: userId,
44+
plan: 'pro',
45+
status: 'active',
46+
periodStart: new Date('2025-02-01T00:00:00.000Z'),
47+
periodEnd: new Date('2025-03-01T00:00:00.000Z'),
48+
}
49+
}
50+
return {
51+
getHighestPrioritySubscription: subscription,
52+
getHighestPriorityPersonalSubscription: subscription,
53+
getOrganizationSubscriptionUsable: vi.fn(),
54+
}
55+
})
56+
vi.mock('@/lib/billing/core/plan', () => ({
57+
getHighestPrioritySubscription: vi.fn(),
58+
getHighestPriorityPersonalSubscription: vi.fn(),
59+
}))
60+
vi.mock('@/lib/billing/core/access', () => ({
61+
getEffectiveBillingStatus: async () => ({ billingBlocked: false }),
62+
isOrganizationBillingBlocked: async () => false,
63+
}))
64+
vi.mock('@/lib/billing/core/billing', () => ({
65+
calculateSubscriptionOverage: async () => 0,
66+
computeOrgOverageAmount: vi.fn(),
67+
getOrganizationSubscription: vi.fn(),
68+
}))
69+
vi.mock('@/lib/billing/cycle-close', () => ({ isSubscriptionCycleCloseCurrent: async () => true }))
70+
vi.mock('@/lib/billing/plan-helpers', () => ({ isEnterprise: () => false, isFree: () => false }))
71+
vi.mock('@/lib/billing/subscriptions/utils', () => ({
72+
hasUsableSubscriptionAccess: () => true,
73+
isOrgScopedSubscription: () => false,
74+
}))
75+
vi.mock('@/lib/billing/calculations/usage-monitor', () => ({
76+
checkBillingBlocked: vi.fn(),
77+
checkBillingEntityBlocked: vi.fn(),
78+
checkOrganizationMemberUsageLimit: vi.fn(),
79+
checkUsageStatus: vi.fn(),
80+
}))
81+
vi.mock('@/lib/billing/webhooks/outbox-handlers', () => ({
82+
OUTBOX_EVENT_TYPES: { STRIPE_THRESHOLD_OVERAGE_INVOICE: 'stripe.threshold-overage-invoice' },
83+
}))
84+
vi.mock('@/lib/core/outbox/service', () => ({ enqueueOutboxEvent: vi.fn() }))
85+
vi.mock('@sim/audit', () => ({ AuditAction: {}, AuditResourceType: {}, recordAudit: vi.fn() }))
86+
vi.mock('@/lib/posthog/server', () => ({ captureServerEvent: vi.fn() }))
87+
vi.mock('@/lib/mothership/request/otel', () => ({
88+
withIncomingGoSpan: (
89+
_headers: unknown,
90+
_name: unknown,
91+
_attrs: unknown,
92+
run: (span: { setAttribute: () => void; setAttributes: () => void }) => unknown
93+
) => run({ setAttribute: vi.fn(), setAttributes: vi.fn() }),
94+
}))
95+
96+
import { POST } from '@/app/api/billing/update-cost/route'
97+
98+
const run = promisify(execFile)
99+
100+
afterAll(async () => {
101+
await state.client?.end()
102+
resetEnvMock()
103+
resetEnvFlagsMock()
104+
})
105+
106+
/**
107+
* Requires an isolated localhost PostgreSQL database and the worker checkout.
108+
* Runs worker receipt accounting/settlement through HTTP into this route.
109+
* Provider calls and subscription fixtures are synthetic; both ledgers and callback are real.
110+
*/
111+
describe.skipIf(!state.databaseUrl || !state.workerDirectory)(
112+
'worker service receipt SQL proof',
113+
() => {
114+
it('settles hosted services, BYOK, media and Agent blocks through the real callback and SQL', async () => {
115+
const databaseUrl = new URL(state.databaseUrl as string)
116+
expect(['127.0.0.1', 'localhost']).toContain(databaseUrl.hostname)
117+
const client = state.client
118+
if (!client) throw new Error('Test database client was not initialized')
119+
setEnv({ INTERNAL_API_SECRET: 'billing-replay-local-secret-0123456789' })
120+
setEnvFlags({ isBillingEnabled: true, isHosted: true })
121+
await client.unsafe(`CREATE SCHEMA "${state.schema}"`)
122+
const statuses: number[] = []
123+
const server = createServer(async (request, response) => {
124+
try {
125+
const chunks: Buffer[] = []
126+
let bytes = 0
127+
for await (const chunk of request) {
128+
const buffer = Buffer.from(chunk)
129+
bytes += buffer.length
130+
if (bytes > 16384) throw new Error('Test request exceeds the callback fixture limit')
131+
chunks.push(buffer)
132+
}
133+
const headers = new Headers()
134+
for (const [key, value] of Object.entries(request.headers)) {
135+
if (typeof value === 'string') headers.set(key, value)
136+
}
137+
const result = await POST(
138+
new NextRequest(`http://127.0.0.1${request.url}`, {
139+
method: 'POST',
140+
headers,
141+
body: Buffer.concat(chunks).toString(),
142+
})
143+
)
144+
statuses.push(result.status)
145+
response.writeHead(result.status, Object.fromEntries(result.headers))
146+
response.end(await result.text())
147+
} catch (error) {
148+
response.writeHead(500)
149+
response.end(String(error))
150+
}
151+
})
152+
try {
153+
await client.unsafe(`CREATE TABLE "user" (id text PRIMARY KEY);
154+
INSERT INTO "user" (id) VALUES ('billing-replay-actor'), ('billing-replay-transient');
155+
CREATE TABLE usage_log (
156+
id text PRIMARY KEY, user_id text NOT NULL, category text NOT NULL, source text NOT NULL,
157+
description text NOT NULL, metadata jsonb, cost numeric NOT NULL, event_key text,
158+
billing_entity_type text, billing_entity_id text, billing_period_start timestamp,
159+
billing_period_end timestamp, workspace_id text, workflow_id text, execution_id text,
160+
created_at timestamp NOT NULL DEFAULT now(),
161+
CONSTRAINT usage_log_user_id_user_id_fk FOREIGN KEY (user_id) REFERENCES "user"(id)
162+
); CREATE UNIQUE INDEX usage_log_event_key_unique ON usage_log(event_key) WHERE event_key IS NOT NULL`)
163+
await new Promise<void>((resolve) => server.listen(0, '127.0.0.1', resolve))
164+
const address = server.address()
165+
if (!address || typeof address === 'string') throw new Error('Expected HTTP server port')
166+
const result = await run('bun', ['test', 'apps/server/test/billing-sim-sql.test.ts'], {
167+
cwd: state.workerDirectory,
168+
env: {
169+
...process.env,
170+
INTERNAL_API_SECRET: 'billing-replay-local-secret-0123456789',
171+
BILLING_WORKER_SIM_URL: `http://127.0.0.1:${address.port}`,
172+
DATABASE_URL: state.databaseUrl,
173+
},
174+
timeout: 60_000,
175+
maxBuffer: 1024 * 1024,
176+
}).catch((error: ExecFileException & { stdout?: string; stderr?: string }) => {
177+
throw new Error([error.message, error.stdout, error.stderr].filter(Boolean).join('\n'), {
178+
cause: error,
179+
})
180+
})
181+
expect(result.stderr).toContain('0 fail')
182+
expect(statuses.every((status) => status === 200 || status === 409)).toBe(true)
183+
const rows =
184+
await client`SELECT source, cost, billing_entity_id, billing_period_start::text AS period_start, billing_period_end::text AS period_end FROM usage_log ORDER BY cost`
185+
expect(rows).toHaveLength(4)
186+
expect(rows.map((row) => Number(row.cost))).toEqual([0.0129, 0.3, 0.3075, 0.3075])
187+
expect(rows[0].source).toBe('mothership_block')
188+
for (const row of rows) {
189+
expect(row.billing_entity_id).toBe('billing-replay-actor')
190+
expect(row.period_start).toBe('2025-02-01 00:00:00')
191+
expect(row.period_end).toBe('2025-03-01 00:00:00')
192+
}
193+
} finally {
194+
try {
195+
if (server.listening) {
196+
await new Promise<void>((resolve, reject) =>
197+
server.close((error) => (error ? reject(error) : resolve()))
198+
)
199+
}
200+
} finally {
201+
await client.unsafe(`DROP SCHEMA "${state.schema}" CASCADE`)
202+
}
203+
}
204+
}, 90_000)
205+
}
206+
)

‎apps/sim/app/api/copilot/byok/route.ts‎

Lines changed: 0 additions & 121 deletions
This file was deleted.

‎apps/sim/app/api/mothership/execute/route.test.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,10 @@ vi.mock('@/lib/mothership/request/session/explicit-abort', () => ({
8383
requestExplicitStreamAbort: mockRequestExplicitStreamAbort,
8484
}))
8585

86+
vi.mock('@/lib/mothership/transport/connection', () => ({
87+
getSimConnection: () => ({ mode: 'checkpoint', channelId: 'f'.repeat(64) }),
88+
}))
89+
8690
vi.mock('@/lib/core/config/env-flags', () => ({
8791
isDocSandboxEnabled: false,
8892
}))

‎apps/sim/app/api/mothership/execute/route.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ import { runHeadlessCopilotLifecycle } from '@/lib/mothership/request/lifecycle/
3636
import { requestExplicitStreamAbort } from '@/lib/mothership/request/session/explicit-abort'
3737
import type { StreamEvent } from '@/lib/mothership/request/types'
3838
import { normalizeSecretMountPolicy } from '@/lib/mothership/secret-mount-policy'
39+
import { getSimConnection } from '@/lib/mothership/transport/connection'
3940
import {
4041
assertActiveWorkspaceAccess,
4142
isWorkspaceAccessDeniedError,
@@ -301,6 +302,7 @@ export const POST = withRouteHandler(async (req: NextRequest) => {
301302
: m
302303
)
303304
const requestPayload: Record<string, unknown> = {
305+
simConnection: getSimConnection(),
304306
messages: wireMessages,
305307
...(useConversationHistory !== undefined ? { useConversationHistory } : {}),
306308
...(modelSelection ? { modelSelection } : {}),

0 commit comments

Comments
 (0)