Skip to content

Commit 08442f4

Browse files
committed
Feed the cloud membership mirror from login and member writes
1 parent aeb1fc1 commit 08442f4

16 files changed

Lines changed: 1250 additions & 266 deletions

‎apps/cloud/package.json‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,8 @@
3232
"db:backfill-org-slugs:dev": "op run --env-file=.env.op -- bun run scripts/backfill-org-slugs.ts",
3333
"db:backfill-subjects:prod": "op run --env-file=.env.production -- bun run scripts/backfill-subjects.ts",
3434
"db:backfill-subjects:dev": "op run --env-file=.env.op -- bun run scripts/backfill-subjects.ts",
35+
"db:backfill-workos-mirror:prod": "op run --env-file=.env.production -- bun run scripts/backfill-workos-mirror.ts",
36+
"db:backfill-workos-mirror:dev": "op run --env-file=.env.op -- bun run scripts/backfill-workos-mirror.ts",
3537
"routes:gen": "bun scripts/gen-routes.ts",
3638
"vendor-wasm": "bun run scripts/vendor-quickjs-wasm.ts"
3739
},
Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
// ---------------------------------------------------------------------------
2+
// One-off data backfill: fill the membership mirror (`accounts` profile
3+
// columns + `memberships` rows, migration 0018) from WorkOS for every
4+
// organization the mirror already knows.
5+
//
6+
// bun run db:backfill-workos-mirror:prod # op run --env-file=.env.production
7+
// bun run db:backfill-workos-mirror:dev # against the local PGlite dev db
8+
//
9+
// For each row in `organizations`: list its active + pending memberships,
10+
// fetch each member's user (concurrency 5), upsert user + membership through
11+
// the same guarded store the request path uses (`auth/workos-mirror-store.ts`).
12+
// Idempotent — the upserts refuse anything older than the stored WorkOS
13+
// `updatedAt`, so re-running is safe and never rewinds a fresher row.
14+
// Pass --dry-run to read and count without writing.
15+
// ---------------------------------------------------------------------------
16+
17+
import { asc } from "drizzle-orm";
18+
import { drizzle } from "drizzle-orm/postgres-js";
19+
import { Effect } from "effect";
20+
import postgres from "postgres";
21+
import { WorkOS } from "@workos-inc/node";
22+
23+
import { backfillWorkOsMirror } from "../src/auth/workos-mirror-backfill";
24+
import { makeWorkOsMirrorStore } from "../src/auth/workos-mirror-store";
25+
import { organizations } from "../src/db/schema";
26+
27+
const dryRun = process.argv.includes("--dry-run");
28+
29+
const connectionString = process.env.DATABASE_URL;
30+
if (!connectionString) {
31+
console.error("DATABASE_URL is not set");
32+
process.exit(1);
33+
}
34+
const apiKey = process.env.WORKOS_API_KEY;
35+
if (!apiKey) {
36+
console.error("WORKOS_API_KEY is not set");
37+
process.exit(1);
38+
}
39+
40+
const usesLocalDatabase =
41+
connectionString.includes("127.0.0.1") || connectionString.includes("localhost");
42+
43+
const sql = postgres(connectionString, {
44+
max: 1,
45+
prepare: false,
46+
...(usesLocalDatabase ? {} : { ssl: "require" as const }),
47+
});
48+
const db = drizzle(sql);
49+
const workos = new WorkOS(apiKey);
50+
51+
// The script boundary: raw SDK / driver promises lifted once, here.
52+
const fromPromise = <A>(fn: () => Promise<A>) =>
53+
Effect.tryPromise({ try: fn, catch: (cause) => cause });
54+
55+
await Effect.runPromise(
56+
backfillWorkOsMirror(
57+
{
58+
listOrganizationIds: () =>
59+
fromPromise(async () => {
60+
const rows = await db
61+
.select({ id: organizations.id })
62+
.from(organizations)
63+
.orderBy(asc(organizations.createdAt));
64+
return rows.map((row) => row.id);
65+
}),
66+
listOrgMembers: (organizationId) =>
67+
fromPromise(async () => {
68+
const page = await workos.userManagement.listOrganizationMemberships({
69+
organizationId,
70+
statuses: ["active", "pending"],
71+
});
72+
return page.listMetadata.after ? page.autoPagination() : page.data;
73+
}),
74+
getUser: (userId) => fromPromise(() => workos.userManagement.getUser(userId)),
75+
},
76+
makeWorkOsMirrorStore(db),
77+
{ dryRun, log: (line) => console.log(line) },
78+
).pipe(Effect.ensuring(Effect.promise(() => sql.end({ timeout: 5 })))),
79+
);

‎apps/cloud/src/account/account-api.ts‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import {
99

1010
import { ApiKeyService } from "../auth/api-keys";
1111
import { UserStoreService } from "../auth/context";
12+
import { WorkOsMirror } from "../auth/workos-mirror";
1213
import { sessionFromSealed, type Session } from "../auth/middleware";
1314
import { WorkOSClient } from "../auth/workos";
1415
import { AutumnService } from "../extensions/billing/service";
@@ -95,10 +96,13 @@ const AccountProviderMiddleware = HttpRouter.middleware<{ provides: AccountProvi
9596
* account service closes over the per-request postgres socket). `AutumnService`
9697
* (the seat-gate) stays a residual requirement, satisfied by the app `boot`.
9798
*/
98-
export const workosAccountMiddleware = (rsLive: Layer.Layer<DbService | UserStoreService>) =>
99-
AccountProviderMiddleware.combine(requestScopedMiddleware(rsLive)).layer;
99+
export const workosAccountMiddleware = (
100+
rsLive: Layer.Layer<DbService | UserStoreService | WorkOsMirror>,
101+
) => AccountProviderMiddleware.combine(requestScopedMiddleware(rsLive)).layer;
100102

101-
export const makeAccountApiLive = (rsLive: Layer.Layer<DbService | UserStoreService>) => {
103+
export const makeAccountApiLive = (
104+
rsLive: Layer.Layer<DbService | UserStoreService | WorkOsMirror>,
105+
) => {
102106
// Cloud builds the WorkOS `AccountProvider` INSIDE the request body (so it
103107
// closes over the per-request postgres socket), so it can't be a self-
104108
// contained `Layer<AccountProvider>` — it combines its own middleware with

‎apps/cloud/src/account/org-api-key-revoke.node.test.ts‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import { ApiKeyService, OrgApiKeyNotFound } from "../auth/api-keys";
88
import { UserStoreService } from "../auth/context";
99
import { ORG_SELECTOR_HEADER } from "../auth/organization";
1010
import { WorkOSClient, type WorkOSClientService } from "../auth/workos";
11+
import { WorkOsMirror } from "../auth/workos-mirror";
1112
import { AutumnService } from "../extensions/billing/service";
1213
import { AccountCaller, workosAccountProvider } from "./workos-account-service";
1314

@@ -119,6 +120,16 @@ const stubUsers = Layer.succeed(UserStoreService)({
119120
),
120121
});
121122

123+
// Revoke changes no membership, so the mirror is never written.
124+
const stubMirror = Layer.succeed(WorkOsMirror)({
125+
upsertUser: () => Effect.die("revoke does not write the membership mirror"),
126+
upsertMembership: () => Effect.die("revoke does not write the membership mirror"),
127+
deleteMembership: () => Effect.die("revoke does not write the membership mirror"),
128+
deleteUser: () => Effect.die("revoke does not write the membership mirror"),
129+
getCursor: () => Effect.die("revoke does not read the events cursor"),
130+
setCursor: () => Effect.die("revoke does not move the events cursor"),
131+
});
132+
122133
const stubAutumn = Layer.succeed(AutumnService)({
123134
use: () => Effect.die("revoke does not touch billing"),
124135
ensureCustomer: () => Effect.die("revoke does not touch billing"),
@@ -156,6 +167,7 @@ const providerWith = (accountId: string) => {
156167
Layer.mergeAll(
157168
stubWorkOS,
158169
stubUsers,
170+
stubMirror,
159171
stubApiKeys,
160172
stubAutumn,
161173
Layer.succeed(AccountCaller)({ session: session(accountId) }),

‎apps/cloud/src/account/workos-account-service.ts‎

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import { ApiKeyService } from "../auth/api-keys";
1212
import { UserStoreService } from "../auth/context";
1313
import type { Session } from "../auth/middleware";
1414
import { WorkOSClient } from "../auth/workos";
15+
import { WorkOsMirror, mirrorMembershipFromWorkOs } from "../auth/workos-mirror";
1516
import { ORG_SELECTOR_HEADER, authorizeOrganizationSelector } from "../auth/organization";
1617
import { AutumnService } from "../extensions/billing/service";
1718
import { forkReportMemberSeats } from "../extensions/billing/member-seats";
@@ -50,7 +51,7 @@ export class AccountCaller extends Context.Service<
5051
// (me / API keys) and `org/handlers.ts` (members / roles / invite / role /
5152
// name). Native WorkOS / store failures are mapped at this boundary onto the
5253
// neutral account errors so the shared UI sees one shape:
53-
// WorkOSError | UserStoreError | ApiKeyManagementError → AccountError
54+
// WorkOSError | UserStoreError | ApiKeyManagementError | WorkOsMirrorError → AccountError
5455
// no organization in session → AccountNoOrganization
5556
// not-an-admin / over-seat-limit / not-allowed → AccountForbidden
5657
// ---------------------------------------------------------------------------
@@ -65,13 +66,17 @@ const toAccountError = () => Effect.fail(new AccountError({ message: "Account re
6566
export const workosAccountProvider: Layer.Layer<
6667
AccountProvider,
6768
never,
68-
WorkOSClient | UserStoreService | ApiKeyService | AutumnService | AccountCaller
69+
WorkOSClient | UserStoreService | WorkOsMirror | ApiKeyService | AutumnService | AccountCaller
6970
> = Layer.effect(AccountProvider)(
7071
Effect.gen(function* () {
7172
const workos = yield* WorkOSClient;
7273
const apiKeys = yield* ApiKeyService;
7374
const autumn = yield* AutumnService;
7475
const users = yield* UserStoreService;
76+
// Membership writes below go to WorkOS FIRST (the authority), then are
77+
// written through to the local mirror so the member list and the seat
78+
// count read the change without waiting for the Events reconciler.
79+
const mirror = yield* WorkOsMirror;
7580

7681
// The caller, resolved once per request by the cookie-only session
7782
// middleware (account-api.ts) — the same credential `SessionAuthLive`
@@ -365,6 +370,9 @@ export const workosAccountProvider: Layer.Layer<
365370
yield* workos
366371
.deleteOrgMembership(membershipId)
367372
.pipe(Effect.catchTag("WorkOSError", toAccountError));
373+
yield* mirror
374+
.deleteMembership(membershipId)
375+
.pipe(Effect.catchTag("WorkOsMirrorError", toAccountError));
368376
yield* forkReportMemberSeats(org.id).pipe(Effect.provideContext(ctx));
369377
return { success: true };
370378
}),
@@ -374,9 +382,12 @@ export const workosAccountProvider: Layer.Layer<
374382
const { session, org } = yield* requireOrganization(headers);
375383
yield* requireAdmin(session.accountId, org.id);
376384
yield* assertMembershipInOrg(org.id, membershipId);
377-
yield* workos
385+
const updated = yield* workos
378386
.updateOrgMembershipRole(membershipId, roleSlug)
379387
.pipe(Effect.catchTag("WorkOSError", toAccountError));
388+
yield* mirror
389+
.upsertMembership(mirrorMembershipFromWorkOs(updated))
390+
.pipe(Effect.catchTag("WorkOsMirrorError", toAccountError));
380391
return { success: true };
381392
}),
382393

‎apps/cloud/src/api/layers.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,9 @@ export const BootSharedServices = Layer.mergeAll(
5656
// `AutumnService.Default` is provided HERE because the `createOrganization`
5757
// handler reads it for the free-organizations-per-user limit gate — one of the
5858
// few app-only billing touchpoints. (It is NOT on the neutral boot core.)
59-
export const makeNonProtectedApiLive = (rsLive: Layer.Layer<DbService | UserStoreService>) =>
59+
export const makeNonProtectedApiLive = (
60+
rsLive: Layer.Layer<DbService | UserStoreService | WorkOsMirror>,
61+
) =>
6062
HttpApiBuilder.layer(NonProtectedApi).pipe(
6163
Layer.provide(Layer.mergeAll(CloudAuthPublicHandlers, CloudSessionAuthHandlers)),
6264
Layer.provide(requestScopedMiddleware(rsLive).layer),

‎apps/cloud/src/api/router.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import { HttpRouter } from "effect/unstable/http";
44
import { RouterConfigLive, requestScopedMiddleware } from "@executor-js/api/server";
55

66
import { UserStoreService } from "../auth/context";
7+
import { WorkOsMirror } from "../auth/workos-mirror";
78
import { DbService } from "../db/db";
89
import { makeAccountApiLive } from "../account/account-api";
910

@@ -29,7 +30,9 @@ import { makeProtectedApiLive } from "./protected";
2930
// so tests can substitute a counting fake for `DbService.Live` and
3031
// assert per-request semantics — see
3132
// `apps/cloud/src/api.request-scope.node.test.ts`.
32-
export const makeApiLive = (requestScopedLive: Layer.Layer<DbService | UserStoreService>) => {
33+
export const makeApiLive = (
34+
requestScopedLive: Layer.Layer<DbService | UserStoreService | WorkOsMirror>,
35+
) => {
3336
const BillingRoutesLive = AutumnRoutesLive.pipe(
3437
Layer.provide(requestScopedMiddleware(requestScopedLive).layer),
3538
);

‎apps/cloud/src/auth/api.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { HttpApiEndpoint, HttpApiGroup } from "effect/unstable/httpapi";
22
import { Schema } from "effect";
3-
import { UserStoreError, WorkOSError } from "./errors";
3+
import { UserStoreError, WorkOSError, WorkOsMirrorError } from "./errors";
44
import { NoOrganization } from "@executor-js/api/server";
55
import { SessionAuth } from "./middleware";
66

@@ -171,7 +171,9 @@ export const AUTH_PATHS = {
171171
callback: "/api/auth/callback",
172172
} as const;
173173

174-
const AuthErrors = [UserStoreError, WorkOSError] as const;
174+
// The login callback and the org handlers feed the membership mirror, so a
175+
// mirror write failure is one of their wire errors (same 500 as a store failure).
176+
const AuthErrors = [UserStoreError, WorkOSError, WorkOsMirrorError] as const;
175177
const McpApprovalErrors = [
176178
NoOrganization,
177179
McpExecutionNotFoundError,

‎apps/cloud/src/auth/errors.ts‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,27 @@ export class UserStoreError extends Schema.TaggedErrorClass<UserStoreError>()(
3838
}
3939
}
4040

41+
/**
42+
* The public failure of every cloud membership-mirror write (`WorkOsMirror`).
43+
* Same two diagnosable fields as `UserStoreError` — which mirror call failed,
44+
* and how — classified from the driver cause the same way. Declared here,
45+
* beside `UserStoreError`, because the auth API (`auth/api.ts`, in the SPA
46+
* bundle) names it on the wire for the login and org handlers that feed the
47+
* mirror; the service itself lives in `workos-mirror.ts`.
48+
*/
49+
export class WorkOsMirrorError extends Schema.TaggedErrorClass<WorkOsMirrorError>()(
50+
"WorkOsMirrorError",
51+
{
52+
operation: Schema.String,
53+
reason: Schema.Literals(USER_STORE_FAILURE_REASONS),
54+
},
55+
{ httpApiStatus: 500 },
56+
) {
57+
override get message(): string {
58+
return `workos mirror ${this.operation} failed: ${this.reason}`;
59+
}
60+
}
61+
4162
/** Reasons a retry can plausibly clear: the query never reached a healthy
4263
* server. A `query` failure is deterministic and must not be retried. */
4364
export const isTransientUserStoreReason = (reason: UserStoreFailureReason): boolean =>

0 commit comments

Comments
 (0)