/** * Test Utilities for Transport Testing * * This module provides common test utilities, mock implementations, * and test data generators for transport testing scenarios. */ import { Effect, Stream, pipe, Chunk, Ref, Duration, Scope } from 'effect'; import { TransportError, makeTransportMessage } from '@codeforbreakfast/eventsourcing-transport'; import type { TransportMessage, ConnectedTransportTestInterface, ConnectionState, } from './test-layer-interfaces'; // ============================================================================ // Test Data Generators // ============================================================================ /** * Generates unique message IDs for testing */ export const generateMessageId = (): string => `msg-${Math.random().toString(36).substring(7)}-${Date.now()}`; /** * Make a test transport message */ export const makeTestTransportMessage = ( type: string, payload: TPayload, options?: { id?: string; metadata?: Record; } ): TransportMessage => ({ id: options?.id || generateMessageId(), type, payload, metadata: options?.metadata, }); // ============================================================================ // Mock Implementations // ============================================================================ /** * Mock transport state for testing */ export interface MockTransportState { readonly messages: readonly TransportMessage[]; readonly subscriptions: readonly ((message: TransportMessage) => void)[]; readonly isConnected: boolean; readonly connectionStateSubscribers: readonly ((state: ConnectionState) => void)[]; } const notifySubscribersOfMessage = (message: TransportMessage, updatedState: MockTransportState) => Effect.sync(() => { updatedState.subscriptions.forEach((subscriber) => subscriber(message)); }); const updateStateWithMessage = (message: TransportMessage) => (s: MockTransportState) => ({ ...s, messages: [...s.messages, message], }); const updateMessagesAndNotifySubscribers = ( message: TransportMessage, stateRef: Ref.Ref ): Effect.Effect => Effect.flatMap(Ref.update(stateRef, updateStateWithMessage(message)), () => Effect.flatMap(Ref.get(stateRef), (updatedState) => notifySubscribersOfMessage(message, updatedState) ) ); const checkConnectedAndPublish = ( state: MockTransportState, message: TransportMessage, stateRef: Ref.Ref ): Effect.Effect => { if (!state.isConnected) { return Effect.fail( new TransportError({ message: 'Transport is not connected', cause: undefined, }) ); } return updateMessagesAndNotifySubscribers(message, stateRef); }; const publishMessage = (stateRef: Ref.Ref, message: TransportMessage) => Effect.flatMap(Ref.get(stateRef), (state) => checkConnectedAndPublish(state, message, stateRef)); const notifyDisconnection = (stateRef: Ref.Ref) => Effect.sync(() => { const state = Effect.runSync(Ref.get(stateRef)); state.connectionStateSubscribers.forEach((subscriber) => subscriber('disconnected' as ConnectionState) ); }); const addTransportFinalizer = ( stateRef: Ref.Ref, transport: ConnectedTransportTestInterface ): Effect.Effect => Effect.map( Effect.addFinalizer(() => notifyDisconnection(stateRef)), () => transport ); const addConnectionStateSubscriber = (subscriber: (state: ConnectionState) => void) => (state: MockTransportState) => ({ ...state, connectionStateSubscribers: [...state.connectionStateSubscribers, subscriber], }); const createTransportFromState = ( stateRef: Ref.Ref ): Effect.Effect => { const connectionStateStream = Stream.async((emit) => { const subscriber = (state: ConnectionState) => { emit(Effect.succeed(Chunk.of(state))); }; Effect.runSync(Ref.update(stateRef, addConnectionStateSubscriber(subscriber))); emit(Effect.succeed(Chunk.of('connected' as ConnectionState))); }); const transport: ConnectedTransportTestInterface = { connectionState: connectionStateStream, publish: (message) => publishMessage(stateRef, message), subscribe: (filter) => Effect.succeed( Stream.async((emit) => { const subscriber = (message: TransportMessage) => { if (!filter || filter(message)) { emit(Effect.succeed(Chunk.of(message))); } }; Effect.runSync( Ref.update(stateRef, (state) => ({ ...state, subscriptions: [...state.subscriptions, subscriber], })) ); return Effect.sync(() => { Effect.runSync( Ref.update(stateRef, (state) => ({ ...state, subscriptions: state.subscriptions.filter((sub) => sub !== subscriber), })) ); }); }) ), }; return addTransportFinalizer(stateRef, transport); }; const makeInitialState = (): Effect.Effect, never, never> => Ref.make({ messages: [], subscriptions: [], isConnected: true, connectionStateSubscribers: [], }); /** * Creates a mock transport for testing */ export const makeMockTransport = (): Effect.Effect< ConnectedTransportTestInterface, never, Scope.Scope > => pipe(makeInitialState(), Effect.flatMap(createTransportFromState)); // ============================================================================ // Test Helpers // ============================================================================ const sleepAndRecheck = ( startTime: number, pollIntervalMs: number, checkCondition: (startTime: number) => Effect.Effect ) => Effect.flatMap(Effect.sleep(Duration.millis(pollIntervalMs)), () => checkCondition(startTime)); const checkConditionOrRetry = ( result: boolean, startTime: number, timeoutMs: number, pollIntervalMs: number, checkCondition: (startTime: number) => Effect.Effect ): Effect.Effect => { if (result) { return Effect.void; } if (Date.now() - startTime >= timeoutMs) { return Effect.fail('timeout' as const); } return sleepAndRecheck(startTime, pollIntervalMs, checkCondition); }; const checkConditionWithRetry = ( startTime: number, condition: () => Effect.Effect, timeoutMs: number, pollIntervalMs: number, checkCondition: (startTime: number) => Effect.Effect ): Effect.Effect => Effect.flatMap(condition(), (result) => checkConditionOrRetry(result, startTime, timeoutMs, pollIntervalMs, checkCondition) ); /** * Waits for a condition to become true within a timeout period */ export const waitForCondition = ( condition: () => Effect.Effect, timeoutMs = 5000, pollIntervalMs = 100 ): Effect.Effect => { const checkCondition = (startTime: number): Effect.Effect => checkConditionWithRetry(startTime, condition, timeoutMs, pollIntervalMs, checkCondition); return checkCondition(Date.now()); }; const flipAndFilterError = ( errorPredicate: (error: E) => boolean, effect: Effect.Effect ): Effect.Effect => Effect.filterOrFail( Effect.flip(effect), errorPredicate, () => new Error('Error did not match predicate') ) as Effect.Effect; /** * Expects an effect to fail with a specific error type */ export const expectError = ( effect: Effect.Effect, errorPredicate: (error: E) => boolean ): Effect.Effect => flipAndFilterError(errorPredicate, effect); const collectWithTimeout = ( timeoutMs: number, stream: Stream.Stream ): Effect.Effect, E | 'timeout', never> => Effect.timeoutFail(Stream.runCollect(stream), { duration: Duration.millis(timeoutMs), onTimeout: () => 'timeout' as const, }) as Effect.Effect, E | 'timeout', never>; /** * Collects all values from a stream within a timeout period */ export const collectStreamWithTimeout = ( stream: Stream.Stream, timeoutMs = 5000 ): Effect.Effect, E | 'timeout', never> => collectWithTimeout(timeoutMs, stream); // ============================================================================ // Client-Server Test Helper Implementations // ============================================================================ const filterTakeAndDrain = ( expectedState: ConnectionState, timeoutMs: number, connectionStateStream: Stream.Stream ): Effect.Effect => Effect.mapError( Effect.timeout( Stream.runDrain( Stream.take( Stream.filter(connectionStateStream, (state) => state === expectedState), 1 ) ), timeoutMs ), () => new Error(`Timeout waiting for connection state: ${expectedState}`) ); /** * Default implementation of waitForConnectionState for testing * Waits for a specific connection state to be emitted */ export const waitForConnectionState = ( connectionStateStream: Stream.Stream, expectedState: ConnectionState, timeoutMs: number = 5000 ): Effect.Effect => filterTakeAndDrain(expectedState, timeoutMs, connectionStateStream); const takeCollectAndTimeout = ( count: number, timeoutMs: number, stream: Stream.Stream ): Effect.Effect => Effect.mapError( Effect.timeout( Effect.map(Stream.runCollect(Stream.take(stream, count)), (chunk) => Array.from(chunk)), timeoutMs ), () => new Error(`Timeout collecting ${count} messages`) ); /** * Default implementation of collectMessages for testing * Collects a specific number of messages from a stream */ export const collectMessages = ( stream: Stream.Stream, count: number, timeoutMs: number = 5000 ): Effect.Effect => takeCollectAndTimeout(count, timeoutMs, stream); /** * Default implementation of makeTestMessage for testing * Creates a test message with unique ID */ export const makeTestMessage = (type: string, payload: unknown): TransportMessage => makeTransportMessage( `test-${Date.now()}-${Math.random()}`, type, typeof payload === 'string' ? payload : JSON.stringify(payload), undefined );