Skip to content

Commit 432d7f5

Browse files
d-csclaude
andcommitted
feat(webapp): build one run-ops client pair per shard descriptor
selectRunOpsTopology gains a shard loop and returns a keyed shard map. An aliasOf:"new" descriptor reuses the new store's clients by reference and opens no pool. Each real shard gets its own resilience budget and the new-role pool knobs merged with its per-shard overrides. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 2222f4c commit 432d7f5

3 files changed

Lines changed: 178 additions & 7 deletions

File tree

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

Lines changed: 89 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ import {
3232
} from "./v3/runOpsMigration/splitMode.server";
3333
import { computeRunOpsSplitReadEnabled } from "./v3/runOpsMigration/runOpsSplitReadGate";
3434
import { resolveRunOpsPoolKnobs } from "./v3/runOpsPoolKnobs.server";
35+
import { resolveShardResilience } from "./v3/transactionResilience.server";
3536
import { assertControlPlaneCoresidencyAdvisory } from "./v3/runOpsMigration/controlPlaneCoresidencySentinel.server";
3637
import { DATASOURCE_CONTEXT_KEY, startActiveSpan } from "./v3/tracer.server";
3738
import {
@@ -276,10 +277,19 @@ export const webhookReplica: WebhookReplicaDatabase = singleton("webhookReplica"
276277

277278
type RunOpsClients = { writer: PrismaClient; replica: PrismaReplicaClient };
278279
type NewRunOpsClients = { writer: RunOpsPrismaClient; replica: RunOpsPrismaClient };
280+
export type ShardTopologyDescriptor = {
281+
key: string;
282+
url?: string;
283+
replicaUrl?: string;
284+
aliasOf?: "new";
285+
};
279286
export type RunOpsTopology = {
280287
newRunOps: NewRunOpsClients;
281288
legacyRunOps: RunOpsClients;
282289
controlPlane: RunOpsClients;
290+
// One client pair per gen-2 shard descriptor. Empty unless RUN_OPS_SHARDS is configured. An
291+
// aliasOf:"new" descriptor maps to the newRunOps pair BY REFERENCE (no new pool).
292+
shards: Map<string, NewRunOpsClients>;
283293
};
284294
export type SelectRunOpsTopologyConfig = {
285295
splitEnabled: boolean;
@@ -289,6 +299,7 @@ export type SelectRunOpsTopologyConfig = {
289299
newReplicaUrl?: string;
290300
// When true, legacy reuses the control-plane client instead of opening its own pool. Defaults to false.
291301
legacySharesControlPlane?: boolean;
302+
shards?: ShardTopologyDescriptor[];
292303
};
293304
export type RunOpsClientBuilders = {
294305
controlPlane: RunOpsClients;
@@ -298,6 +309,10 @@ export type RunOpsClientBuilders = {
298309
// RunOpsPrismaClient double-cast needed): the legacy DB carries the full control-plane schema.
299310
buildLegacyWriter: (url: string, clientType: string) => PrismaClient;
300311
buildLegacyReplica: (url: string, clientType: string) => PrismaReplicaClient;
312+
// Receive the whole descriptor so the singleton can resolve per-shard knobs and resilience by key.
313+
// Optional so the existing test literals (which build no shards) need no change.
314+
buildShardWriter?: (shard: ShardTopologyDescriptor) => RunOpsPrismaClient;
315+
buildShardReplica?: (shard: ShardTopologyDescriptor) => RunOpsPrismaClient;
301316
};
302317

303318
// Pure run-ops client selector. No env, no isSplitEnabled() — those
@@ -316,11 +331,11 @@ export function selectRunOpsTopology(
316331
};
317332

318333
if (!config.splitEnabled) {
319-
return { newRunOps: cpFallback, legacyRunOps: controlPlane, controlPlane };
334+
return { newRunOps: cpFallback, legacyRunOps: controlPlane, controlPlane, shards: new Map() };
320335
}
321336

322337
if (!config.legacyUrl || !config.newUrl) {
323-
return { newRunOps: cpFallback, legacyRunOps: controlPlane, controlPlane };
338+
return { newRunOps: cpFallback, legacyRunOps: controlPlane, controlPlane, shards: new Map() };
324339
}
325340

326341
// Same-DB legacy reuses the control-plane pool; only build a separate pool once the DSNs diverge.
@@ -339,12 +354,28 @@ export function selectRunOpsTopology(
339354
const newReplica: RunOpsPrismaClient = config.newReplicaUrl
340355
? builders.buildNewReplica(config.newReplicaUrl, "run-ops-replica")
341356
: newWriter;
357+
const newRunOps: NewRunOpsClients = { writer: newWriter, replica: newReplica };
358+
359+
const shards = new Map<string, NewRunOpsClients>();
360+
for (const shard of config.shards ?? []) {
361+
if (shard.aliasOf === "new") {
362+
// Aliased: share the new store's clients by reference. No builder, no new pool — the soak path.
363+
shards.set(shard.key, newRunOps);
364+
continue;
365+
}
366+
if (!shard.url || !builders.buildShardWriter || !builders.buildShardReplica) {
367+
throw new Error(
368+
`selectRunOpsTopology: shard "${shard.key}" needs a url and shard builders when not aliased`
369+
);
370+
}
371+
const shardWriter = builders.buildShardWriter(shard);
372+
const shardReplica: RunOpsPrismaClient = shard.replicaUrl
373+
? builders.buildShardReplica(shard)
374+
: shardWriter;
375+
shards.set(shard.key, { writer: shardWriter, replica: shardReplica });
376+
}
342377

343-
return {
344-
newRunOps: { writer: newWriter, replica: newReplica },
345-
legacyRunOps,
346-
controlPlane,
347-
};
378+
return { newRunOps, legacyRunOps, controlPlane, shards };
348379
}
349380

350381
// The env-bound run-ops topology singleton. The split decision uses
@@ -378,6 +409,7 @@ const runOpsTopology: RunOpsTopology = singleton("runOpsTopology", () => {
378409
}
379410

380411
const newPoolKnobs = resolveRunOpsPoolKnobs("new");
412+
const shardDescriptorsByKey = new Map(env.RUN_OPS_SHARDS.map((d) => [d.key, d]));
381413

382414
return selectRunOpsTopology(
383415
{
@@ -387,6 +419,12 @@ const runOpsTopology: RunOpsTopology = singleton("runOpsTopology", () => {
387419
newUrl,
388420
newReplicaUrl: env.RUN_OPS_DATABASE_READ_REPLICA_URL,
389421
legacySharesControlPlane,
422+
shards: env.RUN_OPS_SHARDS.map((d) => ({
423+
key: d.key,
424+
url: d.url,
425+
replicaUrl: d.replicaUrl,
426+
aliasOf: d.aliasOf,
427+
})),
390428
},
391429
{
392430
controlPlane: { writer: prisma, replica: $replica },
@@ -461,6 +499,50 @@ const runOpsTopology: RunOpsTopology = singleton("runOpsTopology", () => {
461499
)
462500
)
463501
),
502+
// A gen-2 shard is a dedicated run-ops DB, so it mirrors buildNewWriter/buildNewReplica: same
503+
// client class, same wrapper stack, its OWN resilience budget, and the "new"-role pool knobs
504+
// merged with the descriptor's per-shard overrides. Shards share the run-ops datasource tag.
505+
buildShardWriter: (shard) => {
506+
const descriptor = shardDescriptorsByKey.get(shard.key);
507+
const knobs = resolveRunOpsPoolKnobs("new", descriptor?.knobs);
508+
return registerTransactionResilience(
509+
captureInfraErrorsRunOps(
510+
tagDatasourceRunOps(
511+
"run-ops-writer",
512+
buildRunOpsClient({
513+
url: shard.url!,
514+
clientType: `run-ops-shard-${shard.key}-writer`,
515+
role: "writer",
516+
connectionLimit: knobs.connectionLimit,
517+
poolTimeout: knobs.writerPoolTimeout,
518+
connectTimeout: knobs.writerConnectionTimeout,
519+
useDriverAdapter: knobs.writerDriverAdapter,
520+
})
521+
)
522+
),
523+
resolveShardResilience(shard.key, descriptor?.knobs)
524+
);
525+
},
526+
buildShardReplica: (shard) => {
527+
const descriptor = shardDescriptorsByKey.get(shard.key);
528+
const knobs = resolveRunOpsPoolKnobs("new", descriptor?.knobs);
529+
return markReadReplicaClient(
530+
captureInfraErrorsRunOps(
531+
tagDatasourceRunOps(
532+
"run-ops-replica",
533+
buildRunOpsClient({
534+
url: shard.replicaUrl!,
535+
clientType: `run-ops-shard-${shard.key}-replica`,
536+
role: "replica",
537+
connectionLimit: knobs.replicaConnectionLimit,
538+
poolTimeout: knobs.replicaPoolTimeout,
539+
connectTimeout: knobs.replicaConnectionTimeout,
540+
useDriverAdapter: knobs.replicaDriverAdapter,
541+
})
542+
)
543+
)
544+
);
545+
},
464546
}
465547
);
466548
});

‎apps/webapp/app/v3/transactionResilience.server.ts‎

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,44 @@ export const runOpsTransactionResilience = resolveTransactionResilience("run-ops
6767
budgetBurst: env.RUN_OPS_DATABASE_TRANSACTION_START_RETRY_BUDGET_BURST,
6868
});
6969

70+
// A gen-2 shard's resilience. Defaults to the RUN_OPS_DATABASE_TRANSACTION_* values (so a shard with
71+
// no overrides matches the gen-1 new store), then applies the descriptor's per-shard overrides. Each
72+
// call builds its OWN budget, so a storm on one shard cannot drain another's.
73+
export function resolveShardResilience(
74+
key: string,
75+
overrides?: {
76+
transactionMaxWaitMs?: number;
77+
transactionStartRetryEnabled?: boolean;
78+
transactionStartRetryMaxAttempts?: number;
79+
transactionStartRetryBackoffMinMs?: number;
80+
transactionStartRetryBackoffMaxMs?: number;
81+
transactionStartRetryBudgetPerSec?: number;
82+
transactionStartRetryBudgetBurst?: number;
83+
}
84+
): TransactionResilienceConfig {
85+
return resolveTransactionResilience(`run-ops-shard-${key}`, {
86+
maxWaitMs: overrides?.transactionMaxWaitMs ?? env.RUN_OPS_DATABASE_TRANSACTION_MAX_WAIT_MS,
87+
enabled:
88+
overrides?.transactionStartRetryEnabled ??
89+
env.RUN_OPS_DATABASE_TRANSACTION_START_RETRY_ENABLED,
90+
maxAttempts:
91+
overrides?.transactionStartRetryMaxAttempts ??
92+
env.RUN_OPS_DATABASE_TRANSACTION_START_RETRY_MAX_ATTEMPTS,
93+
backoffMinMs:
94+
overrides?.transactionStartRetryBackoffMinMs ??
95+
env.RUN_OPS_DATABASE_TRANSACTION_START_RETRY_BACKOFF_MIN_MS,
96+
backoffMaxMs:
97+
overrides?.transactionStartRetryBackoffMaxMs ??
98+
env.RUN_OPS_DATABASE_TRANSACTION_START_RETRY_BACKOFF_MAX_MS,
99+
budgetPerSec:
100+
overrides?.transactionStartRetryBudgetPerSec ??
101+
env.RUN_OPS_DATABASE_TRANSACTION_START_RETRY_BUDGET_PER_SEC,
102+
budgetBurst:
103+
overrides?.transactionStartRetryBudgetBurst ??
104+
env.RUN_OPS_DATABASE_TRANSACTION_START_RETRY_BUDGET_BURST,
105+
});
106+
}
107+
70108
export const runOpsLegacyTransactionResilience = resolveTransactionResilience("run-ops-legacy", {
71109
maxWaitMs: env.RUN_OPS_LEGACY_DATABASE_TRANSACTION_MAX_WAIT_MS,
72110
enabled: env.RUN_OPS_LEGACY_DATABASE_TRANSACTION_START_RETRY_ENABLED,

‎apps/webapp/test/runOpsDbTopology.test.ts‎

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -142,6 +142,57 @@ describe("selectRunOpsTopology (pure)", () => {
142142
expect(topo.legacyRunOps.replica).toBe(legacyWriter);
143143
expect(buildLegacyReplica).not.toHaveBeenCalled();
144144
});
145+
146+
const baseSplit = {
147+
splitEnabled: true,
148+
legacyUrl: "postgres://legacy",
149+
newUrl: "postgres://new",
150+
};
151+
const baseBuilders = () => ({
152+
controlPlane: cp,
153+
buildNewWriter: vi.fn().mockReturnValue({ tag: "nw" } as any),
154+
buildNewReplica: vi.fn().mockReturnValue({ tag: "nr" } as any),
155+
buildLegacyWriter: vi.fn().mockReturnValue({ tag: "lw" } as any),
156+
buildLegacyReplica: vi.fn().mockReturnValue({ tag: "lr" } as any),
157+
});
158+
159+
it("no descriptors: the shards map is empty", () => {
160+
const topo = selectRunOpsTopology(baseSplit, baseBuilders());
161+
expect(topo.shards.size).toBe(0);
162+
});
163+
164+
it("two descriptors: two shard client pairs, each built once", () => {
165+
const buildShardWriter = vi.fn((s: any) => ({ tag: `w:${s.key}` }) as any);
166+
const buildShardReplica = vi.fn((s: any) => ({ tag: `r:${s.key}` }) as any);
167+
const topo = selectRunOpsTopology(
168+
{
169+
...baseSplit,
170+
shards: [
171+
{ key: "a", url: "postgres://a", replicaUrl: "postgres://a-r" },
172+
{ key: "b", url: "postgres://b" },
173+
],
174+
},
175+
{ ...baseBuilders(), buildShardWriter, buildShardReplica }
176+
);
177+
expect(topo.shards.size).toBe(2);
178+
expect(topo.shards.get("a")!.writer).toEqual({ tag: "w:a" });
179+
// b has no replicaUrl, so its replica falls back to its writer (buildShardReplica not called for b).
180+
expect(topo.shards.get("b")!.replica).toEqual({ tag: "w:b" });
181+
expect(buildShardWriter).toHaveBeenCalledTimes(2);
182+
expect(buildShardReplica).toHaveBeenCalledTimes(1);
183+
});
184+
185+
it("an alias descriptor reuses newRunOps by reference and calls no shard builder", () => {
186+
const buildShardWriter = vi.fn();
187+
const buildShardReplica = vi.fn();
188+
const topo = selectRunOpsTopology(
189+
{ ...baseSplit, shards: [{ key: "a", aliasOf: "new" }] },
190+
{ ...baseBuilders(), buildShardWriter, buildShardReplica }
191+
);
192+
expect(topo.shards.get("a")).toBe(topo.newRunOps);
193+
expect(buildShardWriter).not.toHaveBeenCalled();
194+
expect(buildShardReplica).not.toHaveBeenCalled();
195+
});
145196
});
146197

147198
describe("sameDatabaseTarget", () => {

0 commit comments

Comments
 (0)