Skip to content

Commit 84ec665

Browse files
0skiTrigger.dev RepoOps
authored andcommitted
feat(webapp,plugins): provision a metering customer when an organization is created
Adds an opt-in billing plugin, in the same shape as the SSO and RBAC plugins: a contract in `@trigger.dev/plugins`, a loader in `@trigger.dev/billing` that dynamically imports an installed implementation and falls back to a no-op, and the host wiring in the webapp. Mono-RevId: d69c1b3adab90b77f66cf5930ea3516e48a3fe6a
1 parent 1e9372f commit 84ec665

17 files changed

Lines changed: 935 additions & 0 deletions

‎apps/webapp/app/env.server.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2617,6 +2617,10 @@ const EnvironmentSchema = z
26172617
// and emits a `sso.revalidation.timeout` warn log — alert on an
26182618
// elevated rate of those to catch a slow/unhealthy SSO dependency.
26192619
SSO_SESSION_REVALIDATION_TIMEOUT_MS: z.coerce.number().int().positive().default(2000),
2620+
2621+
BILLING_PLUGIN_ENABLED: BoolEnv.default(false),
2622+
BILLING_FORCE_FALLBACK: BoolEnv.default(false),
2623+
BILLING_DATABASE_CONNECTION_LIMIT: z.coerce.number().int().default(2),
26202624
})
26212625
.and(GithubAppEnvSchema)
26222626
.and(S2EnvSchema)

‎apps/webapp/app/models/organization.server.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import {
2727
} from "~/services/platform.v3.server";
2828
import { buildDefaultBillingAlerts } from "~/services/billingAlertsDefaults.server";
2929
import { enqueueAttioWorkspaceSync } from "~/services/attio.server";
30+
import { provisionBillingCustomerForOrg } from "~/services/billingPlugin.server";
3031
import { logger } from "~/services/logger.server";
3132
import { telemetry } from "~/services/telemetry.server";
3233
import {
@@ -153,6 +154,8 @@ export async function createOrganization(
153154
},
154155
});
155156

157+
await provisionBillingCustomerForOrg(organization.id);
158+
156159
// Fire-and-forget; never blocks org creation.
157160
void enqueueAttioWorkspaceSync({
158161
orgId: organization.id,
Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
import billing from "@trigger.dev/billing";
2+
import { prisma } from "~/db.server";
3+
import { env } from "~/env.server";
4+
import { provisionBillingCustomerForNewOrg } from "~/services/provisionBillingCustomer.server";
5+
6+
const billingEnabled = env.BILLING_PLUGIN_ENABLED && !env.BILLING_FORCE_FALLBACK;
7+
8+
const billingPlugin = billing.create({
9+
forceFallback: !billingEnabled,
10+
database: {
11+
writerUrl: env.CONTROL_PLANE_DATABASE_URL ?? env.DATABASE_URL,
12+
writerConnectionLimit: env.BILLING_DATABASE_CONNECTION_LIMIT,
13+
readerConnectionLimit: env.BILLING_DATABASE_CONNECTION_LIMIT,
14+
},
15+
});
16+
17+
export function provisionBillingCustomerForOrg(organizationId: string): Promise<void> {
18+
return provisionBillingCustomerForNewOrg(organizationId, {
19+
enabled: billingEnabled,
20+
controller: billingPlugin,
21+
deleteOrganization: async (id) => {
22+
await prisma.organization.delete({ where: { id } });
23+
},
24+
});
25+
}
Lines changed: 295 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,295 @@
1+
import type {
2+
BillingController,
3+
BillingCustomer,
4+
BillingCustomerError,
5+
ProvisionBillingCustomerParams,
6+
ProvisionBillingCustomerResult,
7+
} from "@trigger.dev/billing";
8+
import { errAsync, okAsync, ResultAsync } from "neverthrow";
9+
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
10+
import {
11+
abortableSleep,
12+
PROVISION_MAX_ATTEMPTS,
13+
provisionBillingCustomerForNewOrg,
14+
type NewOrgProvisionDependencies,
15+
} from "./provisionBillingCustomer.server";
16+
17+
type Outcome = () => ResultAsync<ProvisionBillingCustomerResult, BillingCustomerError>;
18+
19+
const created: Outcome = () =>
20+
okAsync({ organizationId: "org_1", billingCustomerId: "cus_1", outcome: "created" as const });
21+
22+
const never: Outcome = () => new ResultAsync(new Promise(() => {}));
23+
24+
function harness(outcomes: Outcome[], options: { usingPlugin?: () => Promise<boolean> } = {}) {
25+
const calls: ProvisionBillingCustomerParams[] = [];
26+
const deleted: string[] = [];
27+
28+
const controller: BillingController = {
29+
isUsingPlugin: options.usingPlugin ?? (async () => true),
30+
getCustomer(): ResultAsync<BillingCustomer | null, BillingCustomerError> {
31+
return okAsync(null);
32+
},
33+
provisionCustomer(params) {
34+
calls.push(params);
35+
return outcomes[Math.min(calls.length - 1, outcomes.length - 1)]();
36+
},
37+
};
38+
39+
const deps: NewOrgProvisionDependencies = {
40+
enabled: true,
41+
controller,
42+
deleteOrganization: async (id) => {
43+
deleted.push(id);
44+
},
45+
};
46+
47+
return { deps, calls, deleted };
48+
}
49+
50+
function run(deps: NewOrgProvisionDependencies) {
51+
return provisionBillingCustomerForNewOrg("org_1", deps).then(
52+
() => "ok" as const,
53+
(error: Error) => error
54+
);
55+
}
56+
57+
beforeEach(() => {
58+
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout", "performance"] });
59+
});
60+
61+
afterEach(() => {
62+
vi.useRealTimers();
63+
});
64+
65+
describe("provisionBillingCustomerForNewOrg", () => {
66+
it("does nothing when billing is not enabled", async () => {
67+
const { deps, calls, deleted } = harness([created]);
68+
69+
expect(await run({ ...deps, enabled: false })).toBe("ok");
70+
expect(calls).toHaveLength(0);
71+
expect(deleted).toHaveLength(0);
72+
});
73+
74+
it("blocks until the customer is provisioned and keeps the org", async () => {
75+
const { deps, calls, deleted } = harness([created]);
76+
77+
expect(await run(deps)).toBe("ok");
78+
expect(calls).toHaveLength(1);
79+
expect(calls[0].organizationId).toBe("org_1");
80+
expect(calls[0].signal).toBeInstanceOf(AbortSignal);
81+
expect(deleted).toHaveLength(0);
82+
});
83+
84+
it("accepts an org that was already provisioned", async () => {
85+
const { deps, deleted } = harness([
86+
() =>
87+
okAsync({
88+
organizationId: "org_1",
89+
billingCustomerId: "cus_1",
90+
outcome: "already_provisioned" as const,
91+
}),
92+
]);
93+
94+
expect(await run(deps)).toBe("ok");
95+
expect(deleted).toHaveLength(0);
96+
});
97+
98+
it("fails closed when billing is enabled but the plugin did not load", async () => {
99+
const { deps, calls, deleted } = harness([created], { usingPlugin: async () => false });
100+
101+
expect(await run(deps)).toBeInstanceOf(Error);
102+
expect(calls).toHaveLength(0);
103+
expect(deleted).toEqual(["org_1"]);
104+
});
105+
106+
it("keeps the org when the plugin did not load and the policy is skip", async () => {
107+
const { deps, deleted } = harness([created], { usingPlugin: async () => false });
108+
109+
expect(await run({ ...deps, notConfiguredPolicy: "skip" })).toBe("ok");
110+
expect(deleted).toHaveLength(0);
111+
});
112+
113+
it("ends at the deadline even when plugin initialization never settles", async () => {
114+
const { deps, calls, deleted } = harness([created], {
115+
usingPlugin: () => new Promise(() => {}),
116+
});
117+
118+
const settled = run({ ...deps, deadlineMs: 1_000 });
119+
await vi.advanceTimersByTimeAsync(999);
120+
expect(deleted).toHaveLength(0);
121+
await vi.advanceTimersByTimeAsync(1);
122+
123+
expect(await settled).toBeInstanceOf(Error);
124+
expect(calls).toHaveLength(0);
125+
expect(deleted).toEqual(["org_1"]);
126+
});
127+
128+
it("rolls back when checking the plugin throws", async () => {
129+
const { deps, deleted } = harness([created], {
130+
usingPlugin: async () => {
131+
throw new Error("boom");
132+
},
133+
});
134+
135+
expect(await run(deps)).toBeInstanceOf(Error);
136+
expect(deleted).toEqual(["org_1"]);
137+
});
138+
139+
it("retries a transient failure with exponential backoff, then succeeds", async () => {
140+
const { deps, calls, deleted } = harness([
141+
() => errAsync("upstream_unavailable"),
142+
() => errAsync("internal"),
143+
created,
144+
]);
145+
146+
const settled = run(deps);
147+
await vi.advanceTimersByTimeAsync(0);
148+
expect(calls).toHaveLength(1);
149+
await vi.advanceTimersByTimeAsync(199);
150+
expect(calls).toHaveLength(1);
151+
await vi.advanceTimersByTimeAsync(1);
152+
expect(calls).toHaveLength(2);
153+
await vi.advanceTimersByTimeAsync(400);
154+
expect(calls).toHaveLength(3);
155+
156+
expect(await settled).toBe("ok");
157+
expect(deleted).toHaveLength(0);
158+
});
159+
160+
it("retries a provisioning call that throws like a transient failure", async () => {
161+
const { deps, calls } = harness([
162+
() => new ResultAsync(Promise.reject(new Error("socket hang up"))),
163+
created,
164+
]);
165+
166+
const settled = run(deps);
167+
await vi.advanceTimersByTimeAsync(200);
168+
169+
expect(await settled).toBe("ok");
170+
expect(calls).toHaveLength(2);
171+
});
172+
173+
it("retries an in_progress outcome instead of treating it as done", async () => {
174+
const { deps, calls } = harness([
175+
() =>
176+
okAsync({
177+
organizationId: "org_1",
178+
billingCustomerId: null,
179+
outcome: "in_progress" as const,
180+
}),
181+
created,
182+
]);
183+
184+
const settled = run(deps);
185+
await vi.advanceTimersByTimeAsync(200);
186+
187+
expect(await settled).toBe("ok");
188+
expect(calls).toHaveLength(2);
189+
});
190+
191+
it("deletes the org and fails creation once retries are exhausted", async () => {
192+
const { deps, calls, deleted } = harness([() => errAsync("upstream_unavailable")]);
193+
194+
const settled = run(deps);
195+
await vi.advanceTimersByTimeAsync(200 + 400 + 800 + 1600);
196+
const result = await settled;
197+
198+
expect(result).toBeInstanceOf(Error);
199+
expect((result as Error).message).toBe("Organization could not be created.");
200+
expect(calls).toHaveLength(PROVISION_MAX_ATTEMPTS);
201+
expect(deleted).toEqual(["org_1"]);
202+
});
203+
204+
it("fails fast without retrying a permanent failure", async () => {
205+
const { deps, calls, deleted } = harness([() => errAsync("upstream_rejected")]);
206+
207+
expect(await run(deps)).toBeInstanceOf(Error);
208+
expect(calls).toHaveLength(1);
209+
expect(deleted).toEqual(["org_1"]);
210+
});
211+
212+
it("fails org creation when the plugin has no credentials", async () => {
213+
const { deps, calls, deleted } = harness([() => errAsync("not_configured")]);
214+
215+
expect(await run(deps)).toBeInstanceOf(Error);
216+
expect(calls).toHaveLength(1);
217+
expect(deleted).toEqual(["org_1"]);
218+
});
219+
220+
it("keeps the org without a customer when the not-configured policy is skip", async () => {
221+
const { deps, deleted } = harness([() => errAsync("not_configured")]);
222+
223+
expect(await run({ ...deps, notConfiguredPolicy: "skip" })).toBe("ok");
224+
expect(deleted).toHaveLength(0);
225+
});
226+
227+
it("ends at the deadline even when a provisioning call never settles", async () => {
228+
const { deps, calls, deleted } = harness([never]);
229+
230+
const settled = run({ ...deps, deadlineMs: 1_000 });
231+
await vi.advanceTimersByTimeAsync(999);
232+
expect(deleted).toHaveLength(0);
233+
await vi.advanceTimersByTimeAsync(1);
234+
235+
expect(await settled).toBeInstanceOf(Error);
236+
expect(calls).toHaveLength(1);
237+
expect(calls[0].signal?.aborted).toBe(true);
238+
expect(deleted).toEqual(["org_1"]);
239+
});
240+
241+
it("skips a backoff that would run past the deadline", async () => {
242+
const { deps, calls, deleted } = harness([() => errAsync("upstream_unavailable")]);
243+
244+
const settled = run({ ...deps, deadlineMs: 300 });
245+
await vi.advanceTimersByTimeAsync(200);
246+
247+
expect(await settled).toBeInstanceOf(Error);
248+
expect(calls).toHaveLength(2);
249+
expect(deleted).toEqual(["org_1"]);
250+
});
251+
252+
it("still fails creation when the rollback delete itself fails", async () => {
253+
const { deps } = harness([() => errAsync("upstream_rejected")]);
254+
255+
const result = await run({
256+
...deps,
257+
deleteOrganization: async () => {
258+
throw new Error("db down");
259+
},
260+
});
261+
262+
expect((result as Error).message).toBe("Organization could not be created.");
263+
});
264+
});
265+
266+
describe("abortableSleep", () => {
267+
it("resolves as soon as the signal aborts, mid-sleep", async () => {
268+
const controller = new AbortController();
269+
let resolved = false;
270+
void abortableSleep(10_000, controller.signal).then(() => {
271+
resolved = true;
272+
});
273+
274+
await vi.advanceTimersByTimeAsync(100);
275+
expect(resolved).toBe(false);
276+
controller.abort();
277+
await vi.advanceTimersByTimeAsync(0);
278+
279+
expect(resolved).toBe(true);
280+
});
281+
282+
it("resolves after the timeout when the signal never aborts", async () => {
283+
const controller = new AbortController();
284+
let resolved = false;
285+
void abortableSleep(500, controller.signal).then(() => {
286+
resolved = true;
287+
});
288+
289+
await vi.advanceTimersByTimeAsync(499);
290+
expect(resolved).toBe(false);
291+
await vi.advanceTimersByTimeAsync(1);
292+
293+
expect(resolved).toBe(true);
294+
});
295+
});

0 commit comments

Comments
 (0)