Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
154 changes: 146 additions & 8 deletions dev-packages/cloudflare-integration-tests/runner.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
import type { Envelope, EnvelopeItemType } from '@sentry/core';
import type { Envelope, EnvelopeItemType, SerializedStreamedSpan } from '@sentry/core';
import { normalize } from '@sentry/core';
import { createBasicSentryServer } from '@sentry-internal/test-utils';
import { spawn, spawnSync } from 'child_process';
import { existsSync, readdirSync, readFileSync } from 'fs';
import { join } from 'path';
import { inspect } from 'util';
import { expect } from 'vitest';
import { expect, onTestFinished } from 'vitest';
import { getSpansFromEnvelope } from './spanUtils';

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

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

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

/** Either the name of the segment span, or a predicate over it. */
type SegmentMatcher = string | ((segmentSpan: SerializedStreamedSpan) => boolean);

type StartResult = {
completed(): Promise<void>;
/** Every non-ignored envelope received so far, matched or not, for count assertions. */
Expand All @@ -154,6 +158,25 @@ type StartResult = {
expected: Expected | Expected[],
options?: { headers?: Record<string, string>; data?: BodyInit; expectError?: boolean },
): Promise<T | undefined>;
/**
* Accumulates spans across envelopes, grouped by trace, and resolves with the spans of the first
* trace that satisfies `isDone`.
*
* A trace reaches the mock server in more than one envelope: the span buffer flushes on a timer, so
* a segment that is still open when its children flush arrives separately, and a Durable Object or a
* service binding sends its own spans from its own isolate. Anything asserting on a whole trace has
* to accumulate rather than read a single envelope.
*/
collectStreamedSpans(isDone: (spansOfTrace: SerializedStreamedSpan[]) => boolean): Promise<SerializedStreamedSpan[]>;
/**
* Accumulates the spans of a trace until its segment span has arrived.
*
* Only use this to assert on the segment span itself. The segment span ends last, but each
* envelope is its own request to the mock server, so the segment can still be *received* before
* the envelope carrying its children. A suite that asserts on the children has to wait for those
* children by name or by count through `collectStreamedSpans`.
*/
collectStreamedSpansUntilSegment(segment: SegmentMatcher): Promise<SerializedStreamedSpan[]>;
};

/** Creates a test runner */
Expand Down Expand Up @@ -213,17 +236,55 @@ export function createRunner(...paths: string[]) {
return this;
},
start: function (signal?: AbortSignal): StartResult {
const { resolve, reject, promise: isComplete } = deferredPromise(cleanupChildProcesses);
let child: ReturnType<typeof spawn> | undefined;
let childSubWorker: ReturnType<typeof spawn> | undefined;

// Tears down this runner only. `cleanupChildProcesses` tears down every registered runner, so
// running it here would kill a worker another test has already started: a runner whose
// `isComplete` settles after its own test (a suite that asserts on streamed spans never calls
// `completed()`, so the abort signal settles it) would take the next test's worker with it.
// The mock server has to close here as well, otherwise one server per scenario stays listening
// for the whole run.
function cleanupThisRunner(): void {
child?.kill();
childSubWorker?.kill();
closeMockServer?.();
closeMockServer = undefined;
}
Comment on lines +248 to +253

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If cleanupThisRunner runs while createBasicSentryServer is still pending (the test throws after start() but before its first await), closeMockServer and child are both undefined. The .then later opens the server and spawns wrangler dev. Nothing kills it until process exit. This is the leak the PR sets out to fix, on a narrower path.

We could set a disposed flag in cleanupThisRunner. At the top of the .then, and before each spawn, close the server and return if it is set.


// A suite that asserts on streamed spans never calls `completed()`, so `isComplete` never
// settles and its worker would stay alive until the vitest process exits. With one such suite
// per file, a full run ends up with dozens of `wrangler dev` processes competing for the
// machine, and the later suites time out. Tie the teardown to the test instead.
onTestFinished(cleanupThisRunner);

let closeMockServer: (() => void) | undefined;

const { resolve, reject, promise: isComplete } = deferredPromise(cleanupThisRunner);
Comment on lines +259 to +263

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A runner now has three teardown paths:

  1. onTestFinished(cleanupThisRunner) (line 259)
  2. deferredPromise(cleanupThisRunner), when isComplete settles (line 263)
  3. CLEANUP_STEPS, two entries per runner, at process exit (lines 440, 528)

This causes a few problems:

  • With onTestFinished in place, the deferredPromise(cleanupThisRunner) is redundant.
  • CLEANUP_STEPS never drops entries. It grows by two closures per runner and calls kill()/close() on dead objects at exit. It's harmless, but it's extra.
  • The PR description says span-only runners settle "from the abort signal." In vitest 3 the context signal aborts only on timeout, not on normal test end. For a passing span-only test, onTestFinished is the only trigger. The comment at lines 242-247 says the same thing.

Suggestion: Keep onTestFinished. Add cleanupThisRunner to CLEANUP_STEPS once, and have it remove itself when it runs. Drop the deferredPromise(cleanupThisRunner) argument and the two ad hoc CLEANUP_STEPS.add calls. One function, two triggers.


const spanWaiters: {
onSpans: (spans: SerializedStreamedSpan[]) => boolean;
resolve: () => void;
reject: (e: unknown) => void;
}[] = [];
let failure: unknown;

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

const expectedEnvelopeCount = expectedEnvelopes.length;
Expand All @@ -240,8 +301,6 @@ export function createRunner(...paths: string[]) {
workerPortPromise.catch(() => {
// handled in `makeRequest`
});
let child: ReturnType<typeof spawn> | undefined;
let childSubWorker: ReturnType<typeof spawn> | undefined;

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

/** Resolves once `onSpans` returns true for the spans of an arriving span envelope. */
function waitForSpans(onSpans: (spans: SerializedStreamedSpan[]) => boolean): Promise<void> {
return new Promise((resolveWaiter, rejectWaiter) => {
if (failure) {
rejectWaiter(failure);
return;
}
spanWaiters.push({ onSpans, resolve: resolveWaiter, reject: rejectWaiter });
});
}

/**
* Span waiters observe the envelope stream, they never consume from it: a suite can assert on
* streamed spans and on error envelopes at the same time.
*/
function notifySpanWaiters(envelope: Envelope): void {
const spans = getSpansFromEnvelope(envelope);
if (!spans.length) {
return;
}

for (const waiter of spanWaiters.slice()) {
let done: boolean;
try {
done = waiter.onSpans(spans);
} catch (e) {
spanWaiters.splice(spanWaiters.indexOf(waiter), 1);
waiter.reject(e);
continue;
}

if (done) {
spanWaiters.splice(spanWaiters.indexOf(waiter), 1);
waiter.resolve();
}
}
}

function assertEnvelopeMatches(expected: Expected, envelope: Envelope): void {
if (typeof expected === 'function') {
expected(envelope);
Expand All @@ -268,6 +365,8 @@ export function createRunner(...paths: string[]) {
function newEnvelope(envelope: Envelope): void {
if (process.env.DEBUG) log('newEnvelope', inspect(envelope, false, null, true));

notifySpanWaiters(envelope);

const envelopeItemType = envelope[1][0][0].type;

if (ignored.has(envelopeItemType)) {
Expand Down Expand Up @@ -337,6 +436,7 @@ export function createRunner(...paths: string[]) {
createBasicSentryServer(newEnvelope)
.then(async ([mockServerPort, mockServerClose]) => {
if (mockServerClose) {
closeMockServer = mockServerClose;
CLEANUP_STEPS.add(() => {
mockServerClose();
});
Expand Down Expand Up @@ -509,6 +609,44 @@ export function createRunner(...paths: string[]) {
await Promise.all(envelopePromises);
return result;
},
collectStreamedSpans: async function (
isDone: (spansOfTrace: SerializedStreamedSpan[]) => boolean,
): Promise<SerializedStreamedSpan[]> {
const spansByTrace = new Map<string, SerializedStreamedSpan[]>();
let matched: SerializedStreamedSpan[] = [];

await waitForSpans(spans => {
for (const span of spans) {
const spansOfTrace = spansByTrace.get(span.trace_id);
if (spansOfTrace) {
spansOfTrace.push(span);
} else {
spansByTrace.set(span.trace_id, [span]);
}
}

// Every trace is a candidate, so a trace that never satisfies `isDone` cannot hold up the
// one that does. Insertion order means the earliest-arriving trace wins a tie.
for (const spansOfTrace of spansByTrace.values()) {
if (isDone(spansOfTrace)) {
matched = spansOfTrace;
return true;
}
}

return false;
});

return matched;
},
collectStreamedSpansUntilSegment: function (segment: SegmentMatcher): Promise<SerializedStreamedSpan[]> {
const matchesSegment =
typeof segment === 'string' ? (span: SerializedStreamedSpan) => span.name === segment : segment;

return this.collectStreamedSpans(spansOfTrace =>
spansOfTrace.some(span => span.is_segment && matchesSegment(span)),
);
},
};
},
};
Expand Down
18 changes: 18 additions & 0 deletions dev-packages/cloudflare-integration-tests/spanUtils.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
import type { Envelope, SerializedStreamedSpan, SerializedStreamedSpanContainer } from '@sentry/core';

export { getSpanOp } from '@sentry-internal/test-utils';

/**
* The span v2 container of an envelope, or `undefined` when the envelope carries no span item.
*/
export function getSpanContainer(envelope: Envelope): SerializedStreamedSpanContainer | undefined {
const spanItem = envelope[1].find(item => item[0].type === 'span');
return spanItem?.[1] as SerializedStreamedSpanContainer | undefined;
}

/**
* The spans of an envelope, or an empty array when the envelope carries no span item.
*/
export function getSpansFromEnvelope(envelope: Envelope): SerializedStreamedSpan[] {
return getSpanContainer(envelope)?.items ?? [];
}
Loading
Loading