Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { eq, and, desc } from 'drizzle-orm';
import { runEventBus } from '../../../utils/run-events';
import { parseLocation } from '../../../utils/parse-location';
import { validateAndReviveRun } from '../../../utils/revive-run';
import { readShardTokensFromMeta } from '../../../utils/shard-tokens';
import { upsertTraceBlob, findTraceBlob } from '../../../utils/trace-blobs';
import { getStorage } from '../../../storage';
import { joinSuitePath } from '#shared/utils/suites';
Expand Down Expand Up @@ -130,7 +131,12 @@ export default eventHandler(async (event) => {
});
}

await validateAndReviveRun(db, id, testRun, streamToken);
// Accept shard tokens too — in a sharded run every shard uploads its own
// case files, and only one of them holds the run's primary stream token.
const isSharded = !!(testRun.shardTotal && testRun.shardTotal > 1);
const shardTokens = isSharded ? readShardTokensFromMeta(testRun.metadata) : undefined;
const isShardToken = shardTokens ? (token: string) => shardTokens.has(token) : undefined;
await validateAndReviveRun(db, id, testRun, streamToken, isShardToken);

// Locate the run case row the reporter streamed earlier
const { filePath } = parseLocation(caseInfo.location);
Expand Down
55 changes: 43 additions & 12 deletions apps/application/server/api/test-runs/[id]/finish.post.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { sql, eq } from 'drizzle-orm';
import { getDatabase } from '../../../database';
import { testRuns } from '../../../database/schema';
import type { DbClient } from '../../../database';
import { testRuns, testRunsCases } from '../../../database/schema';
import { runEventBus } from '../../../utils/run-events';
import { sanitizeMetadata } from '../../../utils/sanitize';
import { resolveRunBranch } from '../../../utils/run-branch';
Expand All @@ -12,14 +13,45 @@ import { postRunPrFeedbackInBackground } from '../../../utils/scm/pr-feedback';
import { maybeEnqueueHealActionInBackground } from '../../../utils/heal/policy';
import { computeRegressionSignals } from '../../../utils/compute-regression-signals';
import { syncAutoMarkersForRun } from '#shared/handlers/markers';
import { sumFailedAndTimedOut } from '#shared/utils/test-counts';
import { FAILED_STATUS_KEYS, sumFailedAndTimedOut } from '#shared/utils/test-counts';

const FAIL_STATUSES = new Set<string>(FAILED_STATUS_KEYS);

/**
* Whether any test in the run failed, judged by each test's last attempt per browser.
* The `failedTests` counter counts attempts, so it would fail flaky-only runs.
*/
async function hasFinalAttemptFailure(db: DbClient, runId: number): Promise<boolean> {
const rows = await db
.select({
testCaseId: testRunsCases.testCaseId,
browserName: testRunsCases.browserName,
retries: testRunsCases.retries,
status: testRunsCases.status,
})
.from(testRunsCases)
.where(eq(testRunsCases.testRunId, runId));

const finalAttempts = new Map<string, { retries: number; status: string }>();
for (const row of rows) {
const key = `${row.testCaseId}|${row.browserName ?? ''}`;
const retries = row.retries ?? 0;
const prev = finalAttempts.get(key);
if (!prev || retries > prev.retries) finalAttempts.set(key, { retries, status: row.status });
}

for (const attempt of finalAttempts.values()) {
if (FAIL_STATUSES.has(attempt.status)) return true;
}
return false;
}

defineRouteMeta({
openAPI: {
tags: ['Test Runs'],
summary: 'Finish a streaming test run',
description:
'Finalize a streaming test run by setting its final status and calculating performance metrics. Supports pending uploads mode where reports are uploaded asynchronously after finishing. For sharded runs, counters are accumulated and the run finishes only after all shards report.',
'Finalize a streaming test run by setting its final status and calculating performance metrics. Supports pending uploads mode where reports are uploaded asynchronously after finishing. For sharded runs, the run finishes only after all shards report; test counters come from the streamed events.',
parameters: [{ name: 'id', in: 'path', required: true, schema: { type: 'integer' } }],
'x-required-roles': [],
requestBody: {
Expand Down Expand Up @@ -105,9 +137,13 @@ export default eventHandler(async (event) => {
const hasPendingUploads = body.hasPendingUploads === true;

if (isSharded) {
// Sharded run: accumulate counters, track shardsFinished
// Duration: use the maximum across all shards
// Counters: SQL increments to accumulate from multiple shards
// Sharded run: track shardsFinished; duration is the maximum across shards.
// The test counters are NOT touched here — every case (including
// synthesized didnotrun ones) arrives as a streamed event, and the events
// endpoint already increments the counters per inserted row. Adding the
// shards' finish totals on top would count every execution twice. Only
// flakyTests accumulates here: the events tally has no flaky notion, so
// the shards' finish bodies are its single source.

// Merge this shard's durations with any previously accumulated ones
const allDurations: number[] = [];
Expand All @@ -119,12 +155,7 @@ export default eventHandler(async (event) => {
const updateData: Record<string, unknown> = {
updatedAt: new Date(),
status: 'running', // keep running until all shards finish
passedTests: sql`${testRuns.passedTests} + ${body.passedTests ?? 0}`,
failedTests: sql`${testRuns.failedTests} + ${sumFailedAndTimedOut(body.failedTests, body.timedOutTests)}`,
skippedTests: sql`${testRuns.skippedTests} + ${body.skippedTests ?? 0}`,
didNotRunTests: sql`${testRuns.didNotRunTests} + ${body.didNotRunTests ?? 0}`,
flakyTests: sql`${testRuns.flakyTests} + ${flakyTests}`,
totalTests: sql`${testRuns.totalTests} + ${body.totalTests ?? 0}`,
shardsFinished: sql`${testRuns.shardsFinished} + 1`,
// Portable "max of two values": SQLite's scalar MAX(a,b) is an aggregate in
// Postgres, so use a CASE expression that runs on both dialects.
Expand Down Expand Up @@ -153,7 +184,7 @@ export default eventHandler(async (event) => {
updatedRun.shardsFinished >= updatedRun.shardTotal
) {
// All shards done — determine final status
finalStatus = (updatedRun.failedTests ?? 0) > 0 ? 'failed' : 'passed';
finalStatus = (await hasFinalAttemptFailure(db, id)) ? 'failed' : 'passed';

if (allDurations.length > 0) {
const aggStats = durationStats(allDurations);
Expand Down
166 changes: 157 additions & 9 deletions apps/application/tests/sharding.spec.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { createHash } from 'node:crypto';
import { test, expect } from './fixtures';
import { PROJECT } from '#shared/test-project-names';

Expand Down Expand Up @@ -130,6 +131,28 @@ test.describe.serial('Sharding API Tests', () => {
expect((await res.json()).processed).toBe(0);
});

test('shard token is also valid for uploading case files', async ({ request }) => {
// tokenShard1 is not the run's primary stream token — uploads from that
// shard must be accepted via the shard-token fallback.
const traceContent = Buffer.from('Mock shard trace data');
const response = await request.post(`/api/test-runs/${runId}/case-files`, {
multipart: {
streamToken: tokenShard1,
testCase: JSON.stringify({ title: 'shard 1 test C', location: 'tests/shard1.spec.ts:8:3', retries: 1 }),
trace_hash: createHash('sha256').update(traceContent).digest('hex'),
trace: {
name: 'trace.zip',
mimeType: 'application/zip',
buffer: traceContent,
},
},
});
expect(response.ok()).toBeTruthy();
const data = await response.json();
expect(data.success).toBe(true);
expect(data.traces).toBe(1);
});

test('counters reflect events from both shards while still running', async ({ request }) => {
const res = await request.get(`/api/test-runs/${runId}`);
expect(res.ok()).toBeTruthy();
Expand Down Expand Up @@ -166,8 +189,8 @@ test.describe.serial('Sharding API Tests', () => {
const runData = await runRes.json();
expect(runData.status).toBe('running');
expect(runData.shardsFinished).toBe(1);
// Counters should include shard 1's finish totals
expect(runData.failedTests).toBeGreaterThanOrEqual(1);
// Counters reflect the streamed events; finish totals are not added on top
expect(runData.failedTests).toBe(1);
});

test('run finishes with failed status after all shards report', async ({ request }) => {
Expand All @@ -192,18 +215,17 @@ test.describe.serial('Sharding API Tests', () => {
expect(data.status).toBe('failed');
});

test('final counters are properly summed from both shards', async ({ request }) => {
test('final counters count each streamed execution exactly once', async ({ request }) => {
const res = await request.get(`/api/test-runs/${runId}`);
expect(res.ok()).toBeTruthy();
const data = await res.json();
expect(data.status).toBe('failed');

// Events contributed: shard0 2 passed, shard1 1 failed = 3 total, 2 passed, 1 failed
// Finish contributed: shard1 +1 total (+1 failed), shard0 +2 total (+2 passed)
// Final: 3 (events) + 1 + 2 (finish totals) = 6 total
expect(data.totalTests).toBe(6);
expect(data.passedTests).toBe(4); // 2 from events + 2 from shard0 finish
expect(data.failedTests).toBe(2); // 1 from events + 1 from shard1 finish
// Streamed events: shard 0 reported 2 passed, shard 1 reported 1 failed.
// The shards' finish totals must not be added on top of these.
expect(data.totalTests).toBe(3);
expect(data.passedTests).toBe(2);
expect(data.failedTests).toBe(1);
expect(data.shardsFinished).toBe(2);
expect(data.shardTotal).toBe(2);
expect(data.instanceId).toBe(INSTANCE_ID);
Expand Down Expand Up @@ -243,6 +265,132 @@ test.describe.serial('Sharding API Tests', () => {
});
});

test.describe.serial('Sharding: flaky-only run finishes as passed', () => {
const INSTANCE_ID = 'sharding-flaky-only-instance-e2e';
let runId: number;
let tokenShard0: string;
let tokenShard1: string;

test('two shards start and stream a flaky test plus a passing test', async ({ request }) => {
const res0 = await request.post('/api/test-runs/start', {
data: {
projectName: PROJECT.SHARDING_TEST,
startTime: new Date().toISOString(),
instanceId: INSTANCE_ID,
shardIndex: 1,
shardTotal: 2,
},
});
expect(res0.ok()).toBeTruthy();
const data0 = await res0.json();
runId = data0.runId;
tokenShard0 = data0.streamToken;

const res1 = await request.post('/api/test-runs/start', {
data: {
projectName: PROJECT.SHARDING_TEST,
startTime: new Date().toISOString(),
instanceId: INSTANCE_ID,
shardIndex: 2,
shardTotal: 2,
},
});
expect(res1.ok()).toBeTruthy();
tokenShard1 = (await res1.json()).streamToken;

// Shard 0 reports a flaky test: a failed first attempt, then a passed retry.
const events0 = await request.post(`/api/test-runs/${runId}/events`, {
data: {
streamToken: tokenShard0,
testCases: [
{
type: 'complete',
title: 'flaky test',
status: 'failed',
duration: 1200,
location: 'tests/flaky.spec.ts:5:3',
retries: 0,
error: 'Expected element to be visible',
},
{
type: 'complete',
title: 'flaky test',
status: 'passed',
duration: 900,
location: 'tests/flaky.spec.ts:5:3',
retries: 1,
},
],
},
});
expect(events0.ok()).toBeTruthy();
expect((await events0.json()).processed).toBe(2);

// Shard 1 reports a plain passing test
const events1 = await request.post(`/api/test-runs/${runId}/events`, {
data: {
streamToken: tokenShard1,
testCases: [
{
type: 'complete',
title: 'stable test',
status: 'passed',
duration: 700,
location: 'tests/stable.spec.ts:5:3',
retries: 0,
},
],
},
});
expect(events1.ok()).toBeTruthy();
});

test('run finishes as passed although a flaky attempt failed', async ({ request }) => {
// Both shards report 'passed' (Playwright's verdict for a flaky-only run);
// shard 0's counters still carry the failed attempt.
const finish0 = await request.post(`/api/test-runs/${runId}/finish`, {
data: {
streamToken: tokenShard0,
status: 'passed',
duration: 4000,
totalTests: 2,
passedTests: 1,
failedTests: 1,
skippedTests: 0,
flakyTests: 1,
},
});
expect(finish0.ok()).toBeTruthy();
expect((await finish0.json()).status).toBe('running');

const finish1 = await request.post(`/api/test-runs/${runId}/finish`, {
data: {
streamToken: tokenShard1,
status: 'passed',
duration: 3000,
totalTests: 1,
passedTests: 1,
failedTests: 0,
skippedTests: 0,
flakyTests: 0,
},
});
expect(finish1.ok()).toBeTruthy();
// The flaky test's failed attempt must not flip the merged run to failed.
expect((await finish1.json()).status).toBe('passed');

const runRes = await request.get(`/api/test-runs/${runId}`);
expect(runRes.ok()).toBeTruthy();
const runData = await runRes.json();
expect(runData.status).toBe('passed');
expect(runData.flakyTests).toBe(1);
// 3 streamed executions (flaky attempt + retry, stable test), counted once.
expect(runData.totalTests).toBe(3);
expect(runData.passedTests).toBe(2);
expect(runData.failedTests).toBe(1);
});
});

test.describe.serial('Sharding: cross-run instanceId cancellation', () => {
const INSTANCE_A = 'sharding-cancel-instance-a';
const INSTANCE_B = 'sharding-cancel-instance-b';
Expand Down
Loading