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
80 changes: 80 additions & 0 deletions __tests__/e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,8 @@ import {
createServerHandshakeOptions,
} from '../router/handshake';
import { RehandshakeStreamId } from '../transport/message';
import { createPromiseWithResolvers } from '../transport/promises';
import { SessionState } from '../transport/sessionStateMachine';
import { TestSetupHelpers } from '../testUtil/fixtures/transports';

describe.each(testMatrix())(
Expand Down Expand Up @@ -1544,6 +1546,84 @@ describe.each(testMatrix())(
await advanceFakeTimersBySessionGrace();
});

test('a rejected re-handshake after disconnect leaves cleanup to the replacement session', async () => {
const requestSchema = Type.Object({ token: Type.String() });

type ParsedMetadata = Static<typeof requestSchema>;

let token = 'token-v1';
const refreshValidation = createPromiseWithResolvers<
ParsedMetadata | 'REJECTED_BY_CUSTOM_HANDLER'
>();
const refreshValidated = vi.fn();
const construct = vi.fn(() => ({ token }));
const clientTransport = getClientTransport(
'client',
createClientHandshakeOptions(requestSchema, construct),
);
const validate = vi.fn(
async (
metadata: ParsedMetadata,
): Promise<ParsedMetadata | 'REJECTED_BY_CUSTOM_HANDLER'> => {
if (metadata.token === 'token-v1') {
return { token: metadata.token };
}

const result = await refreshValidation.promise;
refreshValidated();

return result;
},
);
const serverTransport = getServerTransport<
typeof requestSchema,
ParsedMetadata
>(
'SERVER',
createServerHandshakeOptions<typeof requestSchema, ParsedMetadata>(
requestSchema,
validate,
),
);
addPostTestCleanup(async () => {
await cleanupTransports([clientTransport, serverTransport]);
});

const protocolError = vi.fn();
serverTransport.addEventListener('protocolError', protocolError);
clientTransport.connect(serverTransport.clientId);
await waitFor(() =>
expect(serverTransport.sessions.get('client')?.state).toBe(
SessionState.Connected,
),
);

clientTransport.reconnectOnConnectionDrop = false;
token = 'token-v2';
expect(serverTransport.requestRehandshake('client')).toBe(true);
await waitFor(() => expect(validate).toHaveBeenCalledTimes(2));

closeAllConnections(clientTransport);
await waitFor(() =>
expect(serverTransport.sessions.get('client')?.state).toBe(
SessionState.NoConnection,
),
);

refreshValidation.resolve('REJECTED_BY_CUSTOM_HANDLER');
await waitFor(() => expect(refreshValidated).toHaveBeenCalledOnce());

expect(protocolError).not.toHaveBeenCalled();
expect(serverTransport.sessions.get('client')?.state).toBe(
SessionState.NoConnection,
);

await advanceFakeTimersBySessionGrace();
await waitFor(() =>
expect(serverTransport.sessions.has('client')).toBe(false),
);
});

test('an in-flight handler observes refreshed metadata mid-stream', async () => {
const requestSchema = Type.Object({ token: Type.String() });

Expand Down
14 changes: 9 additions & 5 deletions transport/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -238,20 +238,24 @@ export abstract class ServerTransport<

/**
* Tears down a session whose re-handshake failed (rejected, malformed, timed
* out, or a thrown validator). No-ops if {@link session} is no longer the live
* session for its peer — a transparent reconnect keeps the same id, so callers
* reaching here after an async gap can't accidentally close the session that
* replaced it.
* out, or a thrown validator). No-ops if {@link session} has been consumed or
* is no longer the live session for its peer — a transparent reconnect keeps
* the same id, so callers reaching here after an async gap can't accidentally
* close the session that replaced it.
*/
private teardownForFailedRehandshake(
session: ServerSession<ConnType>,
reason: string,
) {
if (this.sessions.get(session.to) !== session) {
if (session._isConsumed) {
return;
}

const to = session.to;
if (this.sessions.get(to) !== session) {
return;
}

this.log?.warn(`tearing down session to ${to}: ${reason}`, {
...session.loggingMetadata,
connectedTo: to,
Expand Down
Loading