Skip to content

Commit 1519340

Browse files
committed
test(run-engine,metrics-pipeline): bound emitter-readiness waits and cover the consumer round-trip
An unreachable Redis leaves waitUntilReady pending forever, so the bounded wait fails fast with a descriptive error instead of burning the test timeout. The consumer round-trip test gains the same readiness wait its gauge sibling already had, closing the remaining first-emission drop flake.
1 parent 10941ac commit 1519340

2 files changed

Lines changed: 16 additions & 4 deletions

File tree

internal-packages/metrics-pipeline/src/consumer.test.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ redisTest(
4343
});
4444

4545
await consumer.start();
46+
await emitter.waitUntilReady();
4647
emitter.emit("queueA", { op: "enqueue", q: "queueA" });
4748
emitter.emit("queueB", { op: "started", q: "queueB", wait: 42 });
4849

internal-packages/run-engine/src/run-queue/metrics.test.ts

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,17 @@ const authenticatedEnvDev = {
2424
organization: { id: "o1234" },
2525
};
2626

27+
// A dead Redis leaves waitUntilReady() pending forever (the client retries
28+
// indefinitely), which would burn the whole test timeout with no diagnostic.
29+
async function emitterReady(emitter: MetricsStreamEmitter) {
30+
await Promise.race([
31+
emitter.waitUntilReady(),
32+
setTimeout(15_000).then(() => {
33+
throw new Error("metrics emitter Redis connection never became ready");
34+
}),
35+
]);
36+
}
37+
2738
async function readAllEntries(
2839
redisOptions: {
2940
host: string;
@@ -81,7 +92,7 @@ describe("RunQueue queue-metrics emission", () => {
8192
definition,
8293
flag: { enabled: () => true },
8394
});
84-
await emitter.waitUntilReady();
95+
await emitterReady(emitter);
8596

8697
const queue = new RunQueue({
8798
name: "rq",
@@ -184,7 +195,7 @@ describe("RunQueue queue-metrics emission", () => {
184195
definition,
185196
flag: { enabled: () => true },
186197
});
187-
await emitter.waitUntilReady();
198+
await emitterReady(emitter);
188199
const queue = new RunQueue({
189200
name: "rq",
190201
tracer: trace.getTracer("rq"),
@@ -257,7 +268,7 @@ describe("RunQueue queue-metrics emission", () => {
257268
maxLen: 1000,
258269
};
259270
const emitter = new MetricsStreamEmitter({ redis, definition, flag: { enabled: () => true } });
260-
await emitter.waitUntilReady();
271+
await emitterReady(emitter);
261272
const queue = new RunQueue({
262273
name: "rq",
263274
tracer: trace.getTracer("rq"),
@@ -357,7 +368,7 @@ describe("RunQueue queue-metrics emission", () => {
357368
flag: { enabled: () => true },
358369
gaugeSampleRate: 0,
359370
});
360-
await emitter.waitUntilReady();
371+
await emitterReady(emitter);
361372
const queue = new RunQueue({
362373
name: "rq",
363374
tracer: trace.getTracer("rq"),

0 commit comments

Comments
 (0)