Skip to content

Commit 63b8e6e

Browse files
authored
fix(webapp): scan every run-ops store for the batches list (#4806)
## Summary Batches created on a run-ops store other than the two the list reads were missing from the Batches page. No error, nothing logged: the page just showed fewer batches than exist. This is only reachable once additional run-ops stores are configured, so nothing changes for anyone today. ## Fix The list scanned exactly two databases and merged them by keyset. It now covers one leg per configured store, in ascending precedence order, all issued together. The existing keyset merge generalises without change. Every leg runs the same query, with the same cursor predicate, ordering and over-fetch, so a row's rank within its own leg is never worse than its global rank, and the merged first page is still the true first page. That argument holds for any number of legs, not just two. The empty-state check keeps its existing sequential pair, since a project with no batches is the common case for that path, then issues the remaining checks in a single round trip. A store that declares itself an alias of another shares its client by reference, so it contributes no leg. Scanning it would query the same database twice for rows the other leg already returned. This matches how the routing store and the boot checks treat an alias. The fan-out deliberately fails the page if any store is unreachable, rather than returning a short page. A tolerant merge would recreate the same silent absence this change removes, with a wider blast radius. ## Verification Covered by container tests against real databases: gen-1, legacy and additional stores merged into one ordered page, paging forward and back across a boundary that spans stores, and the empty-state check. Also verified end to end against a live environment with a real corpus: the missing rows reproduce with the new leg removed and appear correctly with it present, ordering interleaves across stores as expected, paging across a store boundary loses and repeats nothing, and the page is byte-identical to before when no extra store is configured. Merge precedence is pinned by its own test: one id seeded on two stores, asserting the higher-authority copy is the one shown. Verified by mutation, since a union-only test passes regardless of leg order. ## Boot interlocks Two related boot checks changed alongside the read path, since configuring an extra store is what makes them reachable. A store configured while split reads are disabled is dropped in silence: no client is built, no leg is added, and rows already resident there disappear from every list with no error. The other two ways the split ends up disabled already refuse to start; this closes the one that did not, and names the stores it is refusing. A store that declares itself an alias of another owns no database, so it is exempt. The distinct-database probe fails closed, which meant one store being briefly unreachable collapsed the deployment to single-DB and then refused the boot entirely. Each target now gets a bounded number of attempts with a short backoff before the probe gives up. Failing closed is unchanged once that budget is exhausted, and a genuine duplicate is still a final answer that is never retried.
1 parent 2e24c01 commit 63b8e6e

10 files changed

Lines changed: 584 additions & 16 deletions

File tree

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

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ import { singleton } from "./utils/singleton";
2828
import { registerDatabaseMetricsSource } from "./utils/databaseMetrics.server";
2929
import {
3030
isSplitEnabled,
31+
assertShardsRequireSplit,
3132
assertSplitRealtimeInterlock,
3233
} from "./v3/runOpsMigration/splitMode.server";
3334
import { computeRunOpsSplitReadEnabled } from "./v3/runOpsMigration/runOpsSplitReadGate";
@@ -617,6 +618,12 @@ export const runOpsSplitReadEnabled: boolean = computeRunOpsSplitReadEnabled({
617618
// interlock). Async, so it cannot live in the synchronous singleton factory — called
618619
// fire-and-forget from the eager-boot path (routing is wired synchronously at module load).
619620
export async function assertRunOpsSplitSentinel(): Promise<void> {
621+
// Shard interlock first: shard clients are only built on the split-on arm, so this case has to be
622+
// checked BEFORE the split-off early return below, which would otherwise skip it in silence.
623+
assertShardsRequireSplit({
624+
splitFlagEnabled: env.RUN_OPS_SPLIT_ENABLED,
625+
shards: env.RUN_OPS_SHARDS,
626+
});
620627
if (!env.RUN_OPS_SPLIT_ENABLED) return;
621628
// Realtime interlock (synchronous): Electric replicates only from the control-plane
622629
// DB, so split-on without the native realtime backend leaves NEW-resident runs

‎apps/webapp/app/presenters/v3/BatchListPresenter.server.ts‎

Lines changed: 26 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,7 @@ export class BatchListPresenter extends BasePresenter {
7575
runOpsNew?: RunOpsPrismaClient; // new run-ops client (run-ops brand ⇒ guard classifies as runops)
7676
runOpsLegacyReplica?: RunOpsPrismaClient; // legacy run-ops READ REPLICA only — never the legacy primary
7777
controlPlaneReplica?: PrismaClientOrTransaction; // control-plane DB (for project)
78+
shardReplicas?: ReadonlyArray<{ key: string; replica: RunOpsPrismaClient }>;
7879
splitEnabled?: boolean; // resolved boot constant
7980
}
8081
) {
@@ -110,9 +111,11 @@ export class BatchListPresenter extends BasePresenter {
110111
// unsound across the residency split: legacy cuid ids ("c…") sort ABOVE new run-ops ids ("0…")
111112
// under id order, so a new-only page can hide pre-flip legacy batches that belong ahead of it.
112113
// Ordering is by createdAt (id tiebreak), which is chronologically correct across both schemes.
113-
const [newRows, legacyRows] = await Promise.all([
114+
const shardReplicas = this.readRoute.shardReplicas ?? [];
115+
const [newRows, legacyRows, ...shardRows] = await Promise.all([
114116
scan(this.readRoute.runOpsNew ?? passthrough),
115117
scan(this.readRoute.runOpsLegacyReplica ?? passthrough),
118+
...shardReplicas.map((shard) => scan(shard.replica)),
116119
]);
117120

118121
// De-dupe by id (new wins), re-sort under the page's keyset order, re-apply the over-fetch LIMIT.
@@ -125,6 +128,11 @@ export class BatchListPresenter extends BasePresenter {
125128
byId.set(row.id, row);
126129
}
127130
}
131+
for (const rows of shardRows) {
132+
for (const row of rows) {
133+
byId.set(row.id, row);
134+
}
135+
}
128136

129137
// forward => newest-first (createdAt DESC), backward => oldest-first (ASC); id is the stable
130138
// tiebreak (ASCII codepoint, NEVER localeCompare).
@@ -167,7 +175,23 @@ export class BatchListPresenter extends BasePresenter {
167175
).batchTaskRun.findFirst({
168176
where: { runtimeEnvironmentId: environmentId },
169177
});
170-
return Boolean(onLegacy);
178+
if (onLegacy) {
179+
return true;
180+
}
181+
182+
const shardReplicas = this.readRoute.shardReplicas ?? [];
183+
if (shardReplicas.length === 0) {
184+
return false;
185+
}
186+
187+
const onShards = await Promise.all(
188+
shardReplicas.map((shard) =>
189+
shard.replica.batchTaskRun.findFirst({
190+
where: { runtimeEnvironmentId: environmentId },
191+
})
192+
)
193+
);
194+
return onShards.some(Boolean);
171195
}
172196

173197
public async call({

‎apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.env.$envParam.batches/route.tsx‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@ import {
5656
runOpsSplitReadEnabled,
5757
type PrismaClientOrTransaction,
5858
} from "~/db.server";
59+
import { runOpsNonAliasedShardReplicas } from "~/v3/runOpsMigration/shardHandles.server";
5960
import {
6061
docsPath,
6162
EnvironmentParamSchema,
@@ -104,6 +105,7 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
104105
runOpsNew: runOpsNewReplicaClient,
105106
runOpsLegacyReplica: runOpsLegacyReplicaClient,
106107
controlPlaneReplica: $replica as unknown as PrismaClientOrTransaction,
108+
shardReplicas: runOpsNonAliasedShardReplicas,
107109
splitEnabled: runOpsSplitReadEnabled,
108110
});
109111
const list = await presenter.call({

‎apps/webapp/app/v3/runOpsMigration/distinctDbSentinel.server.ts‎

Lines changed: 56 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,48 @@ export async function probeControlPlaneCoresidency(
6464

6565
export type DistinctTarget = { id: string; url: string };
6666

67+
/** Injection seam for the retry tests: no containers, no real waiting. */
68+
export type DistinctProbeOptions = {
69+
logger?: { warn: (msg: string, meta?: Record<string, unknown>) => void };
70+
readFingerprint?: (url: string) => Promise<DatabaseFingerprint>;
71+
/** Total attempts per target, including the first. Bounded so boot latency stays bounded. */
72+
attempts?: number;
73+
sleep?: (ms: number) => Promise<void>;
74+
};
75+
76+
const DEFAULT_PROBE_ATTEMPTS = 3;
77+
const RETRY_BASE_DELAY_MS = 250;
78+
79+
const defaultSleep = (ms: number) => new Promise<void>((resolve) => setTimeout(resolve, ms));
80+
81+
/**
82+
* Read one fingerprint, retrying a bounded number of times.
83+
*
84+
* The probe fails CLOSED, and that must not change: "distinct" is a positive claim a failed probe
85+
* cannot support. But failing closed on the first blip means one shard being briefly unreachable
86+
* collapses the deployment to single-DB, and the boot interlock then refuses the boot for the whole
87+
* fleet. A transient error deserves a retry; a persistent one still fails closed, just later.
88+
*/
89+
async function readFingerprintWithRetry(
90+
url: string,
91+
read: (url: string) => Promise<DatabaseFingerprint>,
92+
attempts: number,
93+
sleep: (ms: number) => Promise<void>
94+
): Promise<DatabaseFingerprint> {
95+
let lastError: unknown;
96+
for (let attempt = 1; attempt <= attempts; attempt++) {
97+
try {
98+
return await read(url);
99+
} catch (error) {
100+
lastError = error;
101+
if (attempt < attempts) {
102+
await sleep(RETRY_BASE_DELAY_MS * attempt);
103+
}
104+
}
105+
}
106+
throw lastError;
107+
}
108+
67109
/**
68110
* Set uniqueness over every store that owns its own database. Fail-closed: a probe that cannot
69111
* answer returns NOT distinct, because "distinct" is a positive claim a failed probe cannot support.
@@ -77,14 +119,22 @@ export type DistinctTarget = { id: string; url: string };
77119
*/
78120
export async function probeDistinctStores(
79121
targets: DistinctTarget[],
80-
opts?: { logger?: { warn: (msg: string, meta?: Record<string, unknown>) => void } }
122+
opts?: DistinctProbeOptions
81123
): Promise<{ distinct: true } | { distinct: false; reason: string }> {
82124
if (targets.length < 2) {
83125
return { distinct: true };
84126
}
85127

128+
const read = opts?.readFingerprint ?? readDatabaseFingerprint;
129+
const attempts = opts?.attempts ?? DEFAULT_PROBE_ATTEMPTS;
130+
const sleep = opts?.sleep ?? defaultSleep;
131+
86132
try {
87-
const fingerprints = await Promise.all(targets.map((t) => readDatabaseFingerprint(t.url)));
133+
// Retry per TARGET, not around the whole set: one slow shard must not re-probe the stores that
134+
// already answered. A duplicate verdict below is final and is never retried.
135+
const fingerprints = await Promise.all(
136+
targets.map((t) => readFingerprintWithRetry(t.url, read, attempts, sleep))
137+
);
88138

89139
const seen = new Map<string, string>();
90140
for (const [index, target] of targets.entries()) {
@@ -104,7 +154,9 @@ export async function probeDistinctStores(
104154

105155
return { distinct: true };
106156
} catch (error) {
107-
const reason = `distinct-db sentinel probe failed; failing closed (single-DB). ${String(error)}`;
157+
const reason =
158+
`distinct-db sentinel probe failed after ${opts?.attempts ?? DEFAULT_PROBE_ATTEMPTS} ` +
159+
`attempt(s); failing closed (single-DB). ${String(error)}`;
108160
opts?.logger?.warn(reason, { error });
109161
return { distinct: false, reason };
110162
}
@@ -115,7 +167,7 @@ export async function probeDistinctStores(
115167
export async function probeDistinctDatabases(
116168
legacyUrl: string,
117169
newUrl: string,
118-
opts?: { logger?: { warn: (msg: string, meta?: Record<string, unknown>) => void } }
170+
opts?: DistinctProbeOptions
119171
): Promise<{ distinct: true } | { distinct: false; reason: string }> {
120172
return probeDistinctStores(
121173
[

‎apps/webapp/app/v3/runOpsMigration/shardHandles.server.test.ts‎

Lines changed: 24 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import { describe, expect, it } from "vitest";
2-
import { buildShardHandleMaps } from "./shardHandles.server";
2+
import { buildShardHandleMaps, nonAliasedShardReplicas } from "./shardHandles.server";
33

44
// Two distinct sentinels per shard: the maps must not cross writer and replica.
55
function handle(key: string) {
@@ -35,3 +35,26 @@ describe("buildShardHandleMaps", () => {
3535
expect(replicas.get("a")).not.toEqual({ tag: "a-writer" });
3636
});
3737
});
38+
39+
describe("nonAliasedShardReplicas", () => {
40+
it("yields an empty list when no shard is configured", () => {
41+
expect(nonAliasedShardReplicas([])).toEqual([]);
42+
});
43+
44+
it("keeps the configured order and carries each shard's replica", () => {
45+
expect(nonAliasedShardReplicas([handle("b"), handle("a")])).toEqual([
46+
{ key: "b", replica: { tag: "b-replica" } },
47+
{ key: "a", replica: { tag: "a-replica" } },
48+
]);
49+
});
50+
51+
it("drops a shard that declares aliasOf", () => {
52+
expect(nonAliasedShardReplicas([{ ...handle("a"), aliasOf: "new" }, handle("b")])).toEqual([
53+
{ key: "b", replica: { tag: "b-replica" } },
54+
]);
55+
});
56+
57+
it("never carries a writer in place of a replica", () => {
58+
expect(nonAliasedShardReplicas([handle("a")])[0]?.replica).not.toEqual({ tag: "a-writer" });
59+
});
60+
});

‎apps/webapp/app/v3/runOpsMigration/shardHandles.server.ts‎

Lines changed: 17 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,16 @@
55
* is what keeps every gen-2 arm unreachable today.
66
*/
77
import type { PrismaClient } from "@trigger.dev/database";
8+
import type { RunOpsPrismaClient } from "@internal/run-ops-database";
89
import type { ShardKey } from "@trigger.dev/core/v3/isomorphic";
910
import type { PrismaReplicaClient } from "~/db.server";
1011
import { runOpsShardHandles } from "~/db.server";
1112

1213
type ShardHandle = {
1314
key: string;
14-
writer: unknown;
15-
replica: unknown;
15+
writer: RunOpsPrismaClient;
16+
replica: RunOpsPrismaClient;
17+
aliasOf?: string;
1618
};
1719

1820
export function buildShardHandleMaps(handles: ShardHandle[]): {
@@ -22,8 +24,8 @@ export function buildShardHandleMaps(handles: ShardHandle[]): {
2224
const replicas = new Map<ShardKey, PrismaReplicaClient>();
2325
const writers = new Map<ShardKey, PrismaClient>();
2426
for (const handle of handles) {
25-
replicas.set(handle.key, handle.replica as PrismaReplicaClient);
26-
writers.set(handle.key, handle.writer as PrismaClient);
27+
replicas.set(handle.key, handle.replica as unknown as PrismaReplicaClient);
28+
writers.set(handle.key, handle.writer as unknown as PrismaClient);
2729
}
2830
return { replicas, writers };
2931
}
@@ -40,7 +42,17 @@ function resolveShardHandles(): ShardHandle[] {
4042
}
4143
}
4244

43-
const maps = buildShardHandleMaps(resolveShardHandles());
45+
export function nonAliasedShardReplicas<TClient>(
46+
handles: ReadonlyArray<{ key: string; replica: TClient; aliasOf?: string }>
47+
): ReadonlyArray<{ key: string; replica: TClient }> {
48+
return handles
49+
.filter((handle) => handle.aliasOf === undefined)
50+
.map((handle) => ({ key: handle.key, replica: handle.replica }));
51+
}
52+
53+
const handles = resolveShardHandles();
54+
const maps = buildShardHandleMaps(handles);
4455

4556
export const runOpsShardReplicas = maps.replicas;
4657
export const runOpsShardWriters = maps.writers;
58+
export const runOpsNonAliasedShardReplicas = nonAliasedShardReplicas(handles);

‎apps/webapp/app/v3/runOpsMigration/splitMode.server.ts‎

Lines changed: 34 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,11 @@
77
import { env } from "~/env.server";
88
import { logger } from "~/services/logger.server";
99
import { probeDistinctStores as defaultProbe } from "./distinctDbSentinel.server";
10-
import { nonAliasedShards, type ShardTarget } from "~/v3/runOpsShards.server";
10+
import {
11+
nonAliasedShards,
12+
type RunOpsShardDescriptor,
13+
type ShardTarget,
14+
} from "~/v3/runOpsShards.server";
1115

1216
export type SplitModeConfig = {
1317
flagEnabled: boolean;
@@ -73,6 +77,35 @@ export function assertSplitRealtimeInterlock(config: SplitRealtimeInterlockConfi
7377
}
7478
}
7579

80+
export type ShardsRequireSplitConfig = {
81+
splitFlagEnabled: boolean;
82+
/** Raw descriptors. The alias exemption is applied here so no call site can forget it. */
83+
shards: RunOpsShardDescriptor[];
84+
};
85+
86+
/**
87+
* Boot-time shard interlock (pure predicate). Shard clients are only built on the split-on arm of
88+
* `selectRunOpsTopology`, so a shard configured while the split flag is off is dropped in silence:
89+
* no client, no fan-out leg, and any row already resident on that database vanishes from every
90+
* list with no error. The other two ways split can end up disabled (URLs missing, sentinel not
91+
* distinct) already refuse to boot; this closes the one that does not.
92+
*/
93+
export function assertShardsRequireSplit(config: ShardsRequireSplitConfig): void {
94+
if (config.splitFlagEnabled) {
95+
return;
96+
}
97+
// An aliased shard owns no database: it shares its target's client by reference, so its rows are
98+
// still read with the split off and nothing is dropped. Exempt here exactly as it is exempt from
99+
// the distinctness sentinel, the coresidency loop and replication.
100+
const owning = nonAliasedShards(config.shards).map((shard) => shard.key);
101+
if (owning.length === 0) {
102+
return;
103+
}
104+
throw new Error(
105+
`RUN_OPS_SHARDS configures shard(s) ${owning.join(", ")} but RUN_OPS_SPLIT_ENABLED is off, so no shard client is built and rows on those databases would be silently missing; refusing to start.`
106+
);
107+
}
108+
76109
let cached: Promise<boolean> | undefined;
77110

78111
export function isSplitEnabled(): Promise<boolean> {

0 commit comments

Comments
 (0)