Skip to content

Commit 013f716

Browse files
JPeer264claude
andcommitted
test(cloudflare): Add shared streamed-span helpers to the integration test runner
Adds `spanUtils.ts` next to `runner.ts`, and two span collection methods on the runner. The suites ported to span streaming then do not each write their own envelope handling. `spanUtils.ts` reads a single envelope: `getSpanContainer` and `getSpansFromEnvelope`. It re-exports `getSpanOp` from `@sentry-internal/test-utils`, which this package already depends on. That avoids a third copy of the function in the repo. The runner reads across envelopes: `collectStreamedSpans` and `collectStreamedSpansUntilSegment`. The names follow `dev-packages/test-utils/src/event-proxy-server.ts`, so a streamed span assertion reads the same in this package and in the E2E apps. A trace does not arrive in one envelope when several isolates send spans. The collecting helpers therefore group the spans by trace, and resolve on the first trace that satisfies the predicate. Span waiters observe the envelope stream and never consume it. `.expect(...)` stays usable for error envelopes at the same time. Span waiters also receive the runner rejection, so a worker that fails to boot gives the real error instead of a Vitest timeout. `collectStreamedSpansUntilSegment` is for asserting on the segment span alone. Each envelope is its own request to the mock server, so the segment can be received before the envelope carrying its children even though it ends last. A suite that asserts on the children waits for those children instead, by name or by count. A runner now tears down its own workers from `onTestFinished`, and that teardown kills only its own workers rather than running the shared cleanup. Both halves are needed. A suite that asserts on streamed spans never calls `completed()`, so its runner never settles and its `wrangler dev` would otherwise stay alive until the Vitest process exits. A full run left 32 of them behind. Running the shared cleanup instead would kill the worker the next test had already started. The process-exit cleanup still covers every worker. Renames `public-api/startSpan-streamed` to `public-api/startSpan`, and `tracing/ignoreSpans-streamed` to `tracing/ignoreSpans`. Ports both suites onto the helpers, which proves the shape before the bulk of the port. Span streaming is the default, so a `-streamed` suffix no longer marks a difference, and the explicit `traceLifecycle: 'stream'` those four suites carried says nothing either. It goes with the suffix. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent e57d55c commit 013f716

15 files changed

Lines changed: 464 additions & 333 deletions

File tree

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

Lines changed: 146 additions & 8 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

@@ -139,6 +140,9 @@ function deferredPromise<T = void>(
139140

140141
type Expected = Envelope | ((envelope: Envelope) => void);
141142

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

157180
/** Creates a test runner */
@@ -211,17 +234,55 @@ export function createRunner(...paths: string[]) {
211234
return this;
212235
},
213236
start: function (signal?: AbortSignal): StartResult {
214-
const { resolve, reject, promise: isComplete } = deferredPromise(cleanupChildProcesses);
237+
let child: ReturnType<typeof spawn> | undefined;
238+
let childSubWorker: ReturnType<typeof spawn> | undefined;
239+
240+
// Tears down this runner only. `cleanupChildProcesses` tears down every registered runner, so
241+
// running it here would kill a worker another test has already started: a runner whose
242+
// `isComplete` settles after its own test (a suite that asserts on streamed spans never calls
243+
// `completed()`, so the abort signal settles it) would take the next test's worker with it.
244+
// The mock server has to close here as well, otherwise one server per scenario stays listening
245+
// for the whole run.
246+
function cleanupThisRunner(): void {
247+
child?.kill();
248+
childSubWorker?.kill();
249+
closeMockServer?.();
250+
closeMockServer = undefined;
251+
}
252+
253+
// A suite that asserts on streamed spans never calls `completed()`, so `isComplete` never
254+
// settles and its worker would stay alive until the vitest process exits. With one such suite
255+
// per file, a full run ends up with dozens of `wrangler dev` processes competing for the
256+
// machine, and the later suites time out. Tie the teardown to the test instead.
257+
onTestFinished(cleanupThisRunner);
258+
259+
let closeMockServer: (() => void) | undefined;
260+
261+
const { resolve, reject, promise: isComplete } = deferredPromise(cleanupThisRunner);
262+
263+
const spanWaiters: {
264+
onSpans: (spans: SerializedStreamedSpan[]) => boolean;
265+
resolve: () => void;
266+
reject: (e: unknown) => void;
267+
}[] = [];
268+
let failure: unknown;
215269

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

227288
const expectedEnvelopeCount = expectedEnvelopes.length;
@@ -237,8 +298,6 @@ export function createRunner(...paths: string[]) {
237298
workerPortPromise.catch(() => {
238299
// handled in `makeRequest`
239300
});
240-
let child: ReturnType<typeof spawn> | undefined;
241-
let childSubWorker: ReturnType<typeof spawn> | undefined;
242301

243302
/** Called after each expect callback to check if we're complete */
244303
function expectCallbackCalled(): void {
@@ -254,6 +313,44 @@ export function createRunner(...paths: string[]) {
254313
});
255314
}
256315

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

365+
notifySpanWaiters(envelope);
366+
268367
const envelopeItemType = envelope[1][0][0].type;
269368

270369
if (ignored.has(envelopeItemType)) {
@@ -332,6 +431,7 @@ export function createRunner(...paths: string[]) {
332431
createBasicSentryServer(newEnvelope)
333432
.then(async ([mockServerPort, mockServerClose]) => {
334433
if (mockServerClose) {
434+
closeMockServer = mockServerClose;
335435
CLEANUP_STEPS.add(() => {
336436
mockServerClose();
337437
});
@@ -501,6 +601,44 @@ export function createRunner(...paths: string[]) {
501601
await Promise.all(envelopePromises);
502602
return result;
503603
},
604+
collectStreamedSpans: async function (
605+
isDone: (spansOfTrace: SerializedStreamedSpan[]) => boolean,
606+
): Promise<SerializedStreamedSpan[]> {
607+
const spansByTrace = new Map<string, SerializedStreamedSpan[]>();
608+
let matched: SerializedStreamedSpan[] = [];
609+
610+
await waitForSpans(spans => {
611+
for (const span of spans) {
612+
const spansOfTrace = spansByTrace.get(span.trace_id);
613+
if (spansOfTrace) {
614+
spansOfTrace.push(span);
615+
} else {
616+
spansByTrace.set(span.trace_id, [span]);
617+
}
618+
}
619+
620+
// Every trace is a candidate, so a trace that never satisfies `isDone` cannot hold up the
621+
// one that does. Insertion order means the earliest-arriving trace wins a tie.
622+
for (const spansOfTrace of spansByTrace.values()) {
623+
if (isDone(spansOfTrace)) {
624+
matched = spansOfTrace;
625+
return true;
626+
}
627+
}
628+
629+
return false;
630+
});
631+
632+
return matched;
633+
},
634+
collectStreamedSpansUntilSegment: function (segment: SegmentMatcher): Promise<SerializedStreamedSpan[]> {
635+
const matchesSegment =
636+
typeof segment === 'string' ? (span: SerializedStreamedSpan) => span.name === segment : segment;
637+
638+
return this.collectStreamedSpans(spansOfTrace =>
639+
spansOfTrace.some(span => span.is_segment && matchesSegment(span)),
640+
);
641+
},
504642
};
505643
},
506644
};
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)