Skip to content

Commit 123992d

Browse files
JPeer264claude
andauthored
test(cloudflare): Add shared streamed-span helpers to the integration test runner (#24176)
closes #24146 Adds the shared span assertion helpers `cloudflare-integration-tests` is missing, so the ~80 suites still pinned to `traceLifecycle: 'static'` can be ported without each one hand-rolling its own `getSpanContainer(envelope)`. Three copies of that function exist in the package today. The new helpers are inspired by the E2E tests: - `collectStreamedSpans` - `collectStreamedSpansUntilSegment` We don't need `waitForStreamedSpan` as given in the ticket, as we already have the `.expect` ### Renamed tests - `public-api/startSpan-streamed` is renamed to `public-api/startSpan` - `tracing/ignoreSpans-streamed` to `tracing/ignoreSpans` Both are ported onto the helpers so the shape is proven before the bulk work starts. ### Per-runner teardown A runner now tears down its own worker and its own mock server when it settles, instead of running the shared `cleanupChildProcesses`. A suite that asserts only on streamed spans never calls `completed()`, so its runner settles from the abort signal after its own test has ended. The shared cleanup would then kill the worker the next test had already started, and leaving the mock server open would keep one server per scenario listening for the whole run. The ported suites failed intermittently in large runs until both halves of this were fixed. The process-exit cleanup still covers every runner. --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 5f395ad commit 123992d

15 files changed

Lines changed: 482 additions & 351 deletions

File tree

‎dev-packages/cloudflare-integration-tests/runner.ts‎

Lines changed: 159 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,12 @@
1-
import type { Envelope, EnvelopeItemType } from '@sentry/core';
1+
import type { Envelope, EnvelopeItemType, SerializedStreamedSpan } from '@sentry/core';
22
import { normalize } from '@sentry/core';
33
import { createBasicSentryServer } from '@sentry-internal/test-utils';
44
import { spawn, spawnSync } from 'child_process';
55
import { existsSync, readdirSync, readFileSync } from 'fs';
66
import { join } from 'path';
77
import { inspect } from 'util';
8-
import { expect } from 'vitest';
8+
import { expect, onTestFinished } from 'vitest';
9+
import { getSpansFromEnvelope } from './spanUtils';
910

1011
const CLEANUP_STEPS = new Set<() => void>();
1112

@@ -114,9 +115,11 @@ async function fetchWithRetry(
114115
throw new Error('fetchWithRetry: unreachable');
115116
}
116117

117-
function deferredPromise<T = void>(
118-
done?: () => void,
119-
): { resolve: (val: T) => void; reject: (reason?: unknown) => void; promise: Promise<T> } {
118+
function deferredPromise<T = void>(): {
119+
resolve: (val: T) => void;
120+
reject: (reason?: unknown) => void;
121+
promise: Promise<T>;
122+
} {
120123
let resolve;
121124
let reject;
122125
const promise = new Promise<T>((res, rej) => {
@@ -133,12 +136,15 @@ function deferredPromise<T = void>(
133136
return {
134137
resolve,
135138
reject,
136-
promise: promise.finally(() => done?.()),
139+
promise,
137140
};
138141
}
139142

140143
type Expected = Envelope | ((envelope: Envelope) => void);
141144

145+
/** Either the name of the segment span, or a predicate over it. */
146+
type SegmentMatcher = string | ((segmentSpan: SerializedStreamedSpan) => boolean);
147+
142148
type StartResult = {
143149
completed(): Promise<void>;
144150
/** Every non-ignored envelope received so far, matched or not, for count assertions. */
@@ -154,6 +160,25 @@ type StartResult = {
154160
expected: Expected | Expected[],
155161
options?: { headers?: Record<string, string>; data?: BodyInit; expectError?: boolean },
156162
): Promise<T | undefined>;
163+
/**
164+
* Accumulates spans across envelopes, grouped by trace, and resolves with the spans of the first
165+
* trace that satisfies `isDone`.
166+
*
167+
* A trace reaches the mock server in more than one envelope: the span buffer flushes on a timer, so
168+
* a segment that is still open when its children flush arrives separately, and a Durable Object or a
169+
* service binding sends its own spans from its own isolate. Anything asserting on a whole trace has
170+
* to accumulate rather than read a single envelope.
171+
*/
172+
collectStreamedSpans(isDone: (spansOfTrace: SerializedStreamedSpan[]) => boolean): Promise<SerializedStreamedSpan[]>;
173+
/**
174+
* Accumulates the spans of a trace until its segment span has arrived.
175+
*
176+
* Only use this to assert on the segment span itself. The segment span ends last, but each
177+
* envelope is its own request to the mock server, so the segment can still be *received* before
178+
* the envelope carrying its children. A suite that asserts on the children has to wait for those
179+
* children by name or by count through `collectStreamedSpans`.
180+
*/
181+
collectStreamedSpansUntilSegment(segment: SegmentMatcher): Promise<SerializedStreamedSpan[]>;
157182
};
158183

159184
/** Creates a test runner */
@@ -213,17 +238,54 @@ export function createRunner(...paths: string[]) {
213238
return this;
214239
},
215240
start: function (signal?: AbortSignal): StartResult {
216-
const { resolve, reject, promise: isComplete } = deferredPromise(cleanupChildProcesses);
241+
let child: ReturnType<typeof spawn> | undefined;
242+
let childSubWorker: ReturnType<typeof spawn> | undefined;
243+
let closeMockServer: (() => void) | undefined;
244+
// True after `cleanupThisRunner` ran. The async startup below checks it, so a mock server that
245+
// comes up after the teardown is closed and no `wrangler dev` is spawned.
246+
let disposed = false;
247+
248+
function cleanupThisRunner(): void {
249+
disposed = true;
250+
CLEANUP_STEPS.delete(cleanupThisRunner);
251+
child?.kill();
252+
childSubWorker?.kill();
253+
closeMockServer?.();
254+
closeMockServer = undefined;
255+
}
256+
257+
// A suite that asserts on streamed spans never calls `completed()`, so `isComplete` never
258+
// settles and its worker would stay alive until the vitest process exits. With one such suite
259+
// per file, a full run ends up with dozens of `wrangler dev` processes competing for the
260+
// machine, and the later suites time out. Tie the teardown to the test instead.
261+
onTestFinished(cleanupThisRunner);
262+
CLEANUP_STEPS.add(cleanupThisRunner);
263+
264+
const { resolve, reject, promise: isComplete } = deferredPromise();
265+
266+
const spanWaiters: {
267+
onSpans: (spans: SerializedStreamedSpan[]) => boolean;
268+
resolve: () => void;
269+
reject: (e: unknown) => void;
270+
}[] = [];
271+
let failure: unknown;
217272

218273
// `reject` is called from background event handlers (child process `error`/`exit`, mock server
219274
// callbacks) that fire at arbitrary times relative to the test's `await` points. If `reject` runs
220275
// while nothing is awaiting `isComplete` yet (e.g. a child transiently exits while the test is
221276
// parked in `makeRequest`), the rejection has no handler attached and surfaces as an unhandled
222277
// promise rejection — which Vitest reports as a spurious "Unhandled error" that fails the whole
223-
// suite. Attaching a no-op catch keeps the promise "handled"; the real rejection is still delivered
278+
// suite. Attaching a catch keeps the promise "handled"; the real rejection is still delivered
224279
// to callers via `completed()`, so genuine failures still fail the test.
225-
isComplete.catch(() => {
226-
// handled in `completed()`
280+
//
281+
// A test that only asserts on streamed spans never calls `completed()`, so the same rejection is
282+
// handed to the span waiters as well. Without it, a worker that fails to boot would surface as a
283+
// Vitest timeout instead of the actual error.
284+
isComplete.catch(e => {
285+
failure = e;
286+
for (const waiter of spanWaiters.splice(0)) {
287+
waiter.reject(e);
288+
}
227289
});
228290

229291
const expectedEnvelopeCount = expectedEnvelopes.length;
@@ -240,8 +302,6 @@ export function createRunner(...paths: string[]) {
240302
workerPortPromise.catch(() => {
241303
// handled in `makeRequest`
242304
});
243-
let child: ReturnType<typeof spawn> | undefined;
244-
let childSubWorker: ReturnType<typeof spawn> | undefined;
245305

246306
/** Called after each expect callback to check if we're complete */
247307
function expectCallbackCalled(): void {
@@ -257,6 +317,44 @@ export function createRunner(...paths: string[]) {
257317
});
258318
}
259319

320+
/** Resolves once `onSpans` returns true for the spans of an arriving span envelope. */
321+
function waitForSpans(onSpans: (spans: SerializedStreamedSpan[]) => boolean): Promise<void> {
322+
return new Promise((resolveWaiter, rejectWaiter) => {
323+
if (failure) {
324+
rejectWaiter(failure);
325+
return;
326+
}
327+
spanWaiters.push({ onSpans, resolve: resolveWaiter, reject: rejectWaiter });
328+
});
329+
}
330+
331+
/**
332+
* Span waiters observe the envelope stream, they never consume from it: a suite can assert on
333+
* streamed spans and on error envelopes at the same time.
334+
*/
335+
function notifySpanWaiters(envelope: Envelope): void {
336+
const spans = getSpansFromEnvelope(envelope);
337+
if (!spans.length) {
338+
return;
339+
}
340+
341+
for (const waiter of spanWaiters.slice()) {
342+
let done: boolean;
343+
try {
344+
done = waiter.onSpans(spans);
345+
} catch (e) {
346+
spanWaiters.splice(spanWaiters.indexOf(waiter), 1);
347+
waiter.reject(e);
348+
continue;
349+
}
350+
351+
if (done) {
352+
spanWaiters.splice(spanWaiters.indexOf(waiter), 1);
353+
waiter.resolve();
354+
}
355+
}
356+
}
357+
260358
function assertEnvelopeMatches(expected: Expected, envelope: Envelope): void {
261359
if (typeof expected === 'function') {
262360
expected(envelope);
@@ -268,6 +366,8 @@ export function createRunner(...paths: string[]) {
268366
function newEnvelope(envelope: Envelope): void {
269367
if (process.env.DEBUG) log('newEnvelope', inspect(envelope, false, null, true));
270368

369+
notifySpanWaiters(envelope);
370+
271371
const envelopeItemType = envelope[1][0][0].type;
272372

273373
if (ignored.has(envelopeItemType)) {
@@ -336,11 +436,12 @@ export function createRunner(...paths: string[]) {
336436

337437
createBasicSentryServer(newEnvelope)
338438
.then(async ([mockServerPort, mockServerClose]) => {
339-
if (mockServerClose) {
340-
CLEANUP_STEPS.add(() => {
341-
mockServerClose();
342-
});
439+
// The test can end before the mock server is up, e.g. when it throws right after `start()`.
440+
if (disposed) {
441+
mockServerClose();
442+
return;
343443
}
444+
closeMockServer = mockServerClose;
344445

345446
if (process.env.DEBUG) log('Starting scenario', testPath);
346447

@@ -402,6 +503,10 @@ export function createRunner(...paths: string[]) {
402503
});
403504

404505
await waitForReady(childSubWorker);
506+
507+
if (disposed) {
508+
return;
509+
}
405510
}
406511

407512
child = spawn(
@@ -425,11 +530,6 @@ export function createRunner(...paths: string[]) {
425530
{ stdio: ['ignore', 'pipe', 'inherit'], signal },
426531
);
427532

428-
CLEANUP_STEPS.add(() => {
429-
child?.kill();
430-
childSubWorker?.kill();
431-
});
432-
433533
childSubWorker?.on('error', onChildError);
434534
child.on('error', onChildError);
435535

@@ -509,6 +609,44 @@ export function createRunner(...paths: string[]) {
509609
await Promise.all(envelopePromises);
510610
return result;
511611
},
612+
collectStreamedSpans: async function (
613+
isDone: (spansOfTrace: SerializedStreamedSpan[]) => boolean,
614+
): Promise<SerializedStreamedSpan[]> {
615+
const spansByTrace = new Map<string, SerializedStreamedSpan[]>();
616+
let matched: SerializedStreamedSpan[] = [];
617+
618+
await waitForSpans(spans => {
619+
for (const span of spans) {
620+
const spansOfTrace = spansByTrace.get(span.trace_id);
621+
if (spansOfTrace) {
622+
spansOfTrace.push(span);
623+
} else {
624+
spansByTrace.set(span.trace_id, [span]);
625+
}
626+
}
627+
628+
// Every trace is a candidate, so a trace that never satisfies `isDone` cannot hold up the
629+
// one that does. Insertion order means the earliest-arriving trace wins a tie.
630+
for (const spansOfTrace of spansByTrace.values()) {
631+
if (isDone(spansOfTrace)) {
632+
matched = spansOfTrace;
633+
return true;
634+
}
635+
}
636+
637+
return false;
638+
});
639+
640+
return matched;
641+
},
642+
collectStreamedSpansUntilSegment: function (segment: SegmentMatcher): Promise<SerializedStreamedSpan[]> {
643+
const matchesSegment =
644+
typeof segment === 'string' ? (span: SerializedStreamedSpan) => span.name === segment : segment;
645+
646+
return this.collectStreamedSpans(spansOfTrace =>
647+
spansOfTrace.some(span => span.is_segment && matchesSegment(span)),
648+
);
649+
},
512650
};
513651
},
514652
};
Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
1+
import type { Envelope, SerializedStreamedSpan, SerializedStreamedSpanContainer } from '@sentry/core';
2+
3+
export { getSpanOp } from '@sentry-internal/test-utils';
4+
5+
/**
6+
* The span v2 container of an envelope, or `undefined` when the envelope carries no span item.
7+
*/
8+
export function getSpanContainer(envelope: Envelope): SerializedStreamedSpanContainer | undefined {
9+
const spanItem = envelope[1].find(item => item[0].type === 'span');
10+
return spanItem?.[1] as SerializedStreamedSpanContainer | undefined;
11+
}
12+
13+
/**
14+
* The spans of an envelope, or an empty array when the envelope carries no span item.
15+
*/
16+
export function getSpansFromEnvelope(envelope: Envelope): SerializedStreamedSpan[] {
17+
return getSpanContainer(envelope)?.items ?? [];
18+
}

0 commit comments

Comments
 (0)