import { PushChannelTransfer, SinkTransfer, PollingProxyTransfer, ReadTransfer, LatestStorage, AsyncSinkTransfer, AsyncReadTransfer, AsyncPollingProxyTransfer, linkTransfers, } from '../../src'; import { describe, expect, it, jest } from '@jest/globals'; // ═══════════════════════════════════════════════════════════════ // linkTransfers — Async cases // ═══════════════════════════════════════════════════════════════ // Tests for async linkTransfers strategies (cases 4-9). // Sync cases (1-3) are covered in link-transfers.test.ts. // ═══════════════════════════════════════════════════════════════ // Case 4: Subscribable → AsyncPushable (Reactive stream + async-push) // ═══════════════════════════════════════════════════════════════ describe.each([ ...dataProviderForSubscribableToAsyncPushable(), ] as Array<[number]>)( 'linkTransfers connects Subscribable source to AsyncPushable target test', (value: number) => { it('', async () => { const source = new PushChannelTransfer(); const received: number[] = []; const target = new AsyncSinkTransfer({ callback: (v) => { received.push(v); } }); expect(source.isSubscribable).toBe(true); expect(target.isAsyncPushable).toBe(true); expect(target.isPushable).toBe(false); const subscriber = linkTransfers(source, target); expect(subscriber.active).toBe(true); source.push(value); // asyncPush executes asynchronously — waiting for microtask await new Promise((resolve) => setTimeout(resolve, 10)); expect(received).toEqual([value]); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); /** * Data provider for Subscribable → AsyncPushable. */ function dataProviderForSubscribableToAsyncPushable(): Array { return [ [1], [42], [-7], ]; } describe( 'linkTransfers Subscribable to AsyncPushable forwards multiple values test', () => { it('', async () => { const source = new PushChannelTransfer(); const received: number[] = []; const target = new AsyncSinkTransfer({ callback: (v) => { received.push(v); } }); linkTransfers(source, target); source.push(1); source.push(2); source.push(3); await new Promise((resolve) => setTimeout(resolve, 10)); expect(received).toEqual([1, 2, 3]); source.destroy(); target.destroy(); }); }, ); describe( 'linkTransfers Subscribable to AsyncPushable unsubscribe stops data flow test', () => { it('', async () => { const source = new PushChannelTransfer(); const received: number[] = []; const target = new AsyncSinkTransfer({ callback: (v) => { received.push(v); } }); const subscriber = linkTransfers(source, target); source.push(1); await new Promise((resolve) => setTimeout(resolve, 10)); subscriber.unsubscribe(); expect(subscriber.active).toBe(false); source.push(2); await new Promise((resolve) => setTimeout(resolve, 10)); expect(received).toEqual([1]); source.destroy(); target.destroy(); }); }, ); describe( 'linkTransfers Subscribable to AsyncPushable source not disrupted by target rejection test', () => { it('', async () => { const source = new PushChannelTransfer(); // AsyncSinkTransfer with onError suppresses at the target level, // so asyncPush resolves and .catch() in linkTransfers never fires. // The test verifies that the source subscription remains active // even when the target's callback throws. // // Without link onError: handleError rethrows inside .catch(), // producing an unhandled promise rejection (tested in handle-error.test.ts). // Jest 30 intercepts unhandled rejections via V8 promise hooks, // not process.on('unhandledRejection'), so it cannot be captured here. const target = new AsyncSinkTransfer({ callback: async () => { throw new Error('push error'); }, onError: () => {}, }); const subscriber = linkTransfers(source, target); expect(subscriber.active).toBe(true); source.push(42); await new Promise((resolve) => setTimeout(resolve, 10)); // Subscription remains active — source not disrupted by downstream rejection expect(subscriber.active).toBe(true); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); describe( 'linkTransfers Subscribable to AsyncPushable rejection calls onError test', () => { it('', async () => { const source = new PushChannelTransfer(); const target = new AsyncSinkTransfer({ callback: async () => { throw new Error('push error'); }, }); const onError = jest.fn(); const subscriber = linkTransfers(source, target, { onError }); source.push(42); await new Promise((resolve) => setTimeout(resolve, 10)); expect(onError).toHaveBeenCalledTimes(1); expect(onError).toHaveBeenCalledWith(expect.any(Error), target); expect((onError.mock.calls[0][0] as Error).message).toBe('push error'); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); describe( 'linkTransfers Subscribable to AsyncPushable non-Error rejection wrapped to Error test', () => { it('', async () => { const source = new PushChannelTransfer(); const target = new AsyncSinkTransfer({ callback: async () => { throw 'string error'; }, }); const onError = jest.fn(); const subscriber = linkTransfers(source, target, { onError }); source.push(42); await new Promise((resolve) => setTimeout(resolve, 10)); expect(onError).toHaveBeenCalledTimes(1); expect(onError).toHaveBeenCalledWith(expect.any(Error), target); expect((onError.mock.calls[0][0] as Error).message).toBe('string error'); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); describe( 'linkTransfers Subscribable to AsyncPushable non-Error rejection from raw asyncPush test', () => { it('', async () => { const source = new PushChannelTransfer(); // Mock async pushable that rejects with a non-Error value directly, // bypassing handleError wrapping in AsyncSinkTransfer. const target = { isInput: true, isOutput: false, isDuplex: false, isPushable: false, isPullable: false, isSubscribable: false, isPollingSource: false, isPollingProxy: false, isTriggerable: false, isGate: false, isAsyncPushable: true, isAsyncPullable: false, isAsyncTriggerable: false, isAsyncPollingProxy: false, asyncPush: async () => { throw 'raw string error'; }, destroy: () => {}, } as any; const onError = jest.fn(); const subscriber = linkTransfers(source, target, { onError }); source.push(42); await new Promise((resolve) => setTimeout(resolve, 10)); expect(onError).toHaveBeenCalledTimes(1); expect(onError).toHaveBeenCalledWith(expect.any(Error), target); expect((onError.mock.calls[0][0] as Error).message).toBe('raw string error'); subscriber.unsubscribe(); source.destroy(); }); }, ); // ═══════════════════════════════════════════════════════════════ // Case 5: AsyncPullable → AsyncPollingProxy (Active async-polling) // ═══════════════════════════════════════════════════════════════ describe.each([ ...dataProviderForAsyncPullableToAsyncPollingProxy(), ] as Array<[number]>)( 'linkTransfers connects AsyncPullable source to AsyncPollingProxy target test', (value: number) => { it('', async () => { const mockFlow = { read: jest.fn(async () => value) }; const source = new AsyncReadTransfer({ flow: mockFlow }); const target = new AsyncPollingProxyTransfer({ interval: 10, activated: true }); expect(source.isAsyncPullable).toBe(true); expect(source.isPullable).toBe(false); expect(source.isSubscribable).toBe(false); expect(target.isAsyncPollingProxy).toBe(true); expect(target.isPollingProxy).toBe(false); const received: (number | undefined)[] = []; target.subscribe((data) => { if (data !== undefined) received.push(data); }); const subscriber = linkTransfers(source, target); expect(subscriber.active).toBe(true); await new Promise((resolve) => setTimeout(resolve, 50)); expect(received.length).toBeGreaterThan(0); expect(received[0]).toBe(value); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); /** * Data provider for AsyncPullable → AsyncPollingProxy. */ function dataProviderForAsyncPullableToAsyncPollingProxy(): Array { return [ [1], [42], [-7], ]; } describe( 'linkTransfers AsyncPullable to AsyncPollingProxy unsubscribe stops polling test', () => { it('', async () => { const mockFlow = { read: jest.fn(async () => 42) }; const source = new AsyncReadTransfer({ flow: mockFlow }); const target = new AsyncPollingProxyTransfer({ interval: 10, activated: true }); const received: (number | undefined)[] = []; target.subscribe((data) => { if (data !== undefined) received.push(data); }); const subscriber = linkTransfers(source, target); await new Promise((resolve) => setTimeout(resolve, 50)); const countBefore = received.length; expect(countBefore).toBeGreaterThan(0); subscriber.unsubscribe(); await new Promise((resolve) => setTimeout(resolve, 50)); const countAfter = received.length; expect(countAfter).toBe(countBefore); source.destroy(); target.destroy(); }); }, ); describe( 'linkTransfers AsyncPullable to AsyncPollingProxy subscriber inactive after unsubscribe test', () => { it('', () => { const mockFlow = { read: jest.fn(async () => 1) }; const source = new AsyncReadTransfer({ flow: mockFlow }); const target = new AsyncPollingProxyTransfer({ interval: 100, activated: false }); const subscriber = linkTransfers(source, target); expect(subscriber.active).toBe(true); subscriber.unsubscribe(); expect(subscriber.active).toBe(false); source.destroy(); target.destroy(); }); }, ); // ═══════════════════════════════════════════════════════════════ // Case 6: Pullable → AsyncPollingProxy (Sync-pull wrapped in an async-fetcher) // ═══════════════════════════════════════════════════════════════ describe.each([ ...dataProviderForPullableToAsyncPollingProxy(), ] as Array<[number]>)( 'linkTransfers connects Pullable source to AsyncPollingProxy target test', (value: number) => { it('', async () => { const storage = new LatestStorage(); storage.write(value); const source = new ReadTransfer({ flow: storage }); const target = new AsyncPollingProxyTransfer({ interval: 10, activated: true }); expect(source.isPullable).toBe(true); expect(source.isSubscribable).toBe(false); expect(source.isAsyncPullable).toBe(false); expect(target.isAsyncPollingProxy).toBe(true); const received: (number | undefined)[] = []; target.subscribe((data) => { if (data !== undefined) received.push(data); }); const subscriber = linkTransfers(source, target); expect(subscriber.active).toBe(true); await new Promise((resolve) => setTimeout(resolve, 50)); expect(received.length).toBeGreaterThan(0); expect(received[0]).toBe(value); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); /** * Data provider for Pullable → AsyncPollingProxy. */ function dataProviderForPullableToAsyncPollingProxy(): Array { return [ [1], [42], [-7], ]; } describe( 'linkTransfers Pullable to AsyncPollingProxy unsubscribe stops polling test', () => { it('', async () => { const storage = new LatestStorage(); storage.write(42); const source = new ReadTransfer({ flow: storage }); const target = new AsyncPollingProxyTransfer({ interval: 10, activated: true }); const received: (number | undefined)[] = []; target.subscribe((data) => { if (data !== undefined) received.push(data); }); const subscriber = linkTransfers(source, target); await new Promise((resolve) => setTimeout(resolve, 50)); const countBefore = received.length; expect(countBefore).toBeGreaterThan(0); subscriber.unsubscribe(); await new Promise((resolve) => setTimeout(resolve, 50)); expect(received.length).toBe(countBefore); source.destroy(); target.destroy(); }); }, ); // ═══════════════════════════════════════════════════════════════ // Case 7: Subscribable → AsyncPollingProxy (Subscription + buffer + async-fetcher) // ═══════════════════════════════════════════════════════════════ describe.each([ ...dataProviderForSubscribableToAsyncPollingProxy(), ] as Array<[number]>)( 'linkTransfers connects Subscribable source to AsyncPollingProxy target test', (value: number) => { it('', async () => { const source = new PushChannelTransfer(); const target = new AsyncPollingProxyTransfer({ interval: 10, activated: true }); expect(source.isSubscribable).toBe(true); expect(source.isPullable).toBe(false); expect(target.isAsyncPollingProxy).toBe(true); expect(target.isPushable).toBe(false); expect(target.isPollingProxy).toBe(false); const received: (number | undefined)[] = []; target.subscribe((data) => { if (data !== undefined) received.push(data); }); const subscriber = linkTransfers(source, target); expect(subscriber.active).toBe(true); source.push(value); await new Promise((resolve) => setTimeout(resolve, 50)); expect(received.length).toBeGreaterThan(0); expect(received).toContain(value); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); /** * Data provider for Subscribable → AsyncPollingProxy. */ function dataProviderForSubscribableToAsyncPollingProxy(): Array { return [ [1], [42], [-7], ]; } describe( 'linkTransfers Subscribable to AsyncPollingProxy unsubscribe stops data flow test', () => { it('', async () => { const source = new PushChannelTransfer(); const target = new AsyncPollingProxyTransfer({ interval: 10, activated: true }); const received: (number | undefined)[] = []; target.subscribe((data) => { if (data !== undefined) received.push(data); }); const subscriber = linkTransfers(source, target); source.push(1); await new Promise((resolve) => setTimeout(resolve, 50)); const countBefore = received.length; expect(countBefore).toBeGreaterThan(0); subscriber.unsubscribe(); source.push(2); await new Promise((resolve) => setTimeout(resolve, 50)); expect(received.length).toBe(countBefore); source.destroy(); target.destroy(); }); }, ); describe( 'linkTransfers Subscribable to AsyncPollingProxy buffers last value test', () => { it('', async () => { const source = new PushChannelTransfer(); const target = new AsyncPollingProxyTransfer({ interval: 20, activated: false }); const received: (number | undefined)[] = []; target.subscribe((data) => { if (data !== undefined) received.push(data); }); const subscriber = linkTransfers(source, target); // Sending values before activating polling source.push(10); source.push(20); source.push(30); // Activating polling — the poller will pick up the last value (30) target.activate(); await new Promise((resolve) => setTimeout(resolve, 60)); expect(received).toContain(30); subscriber.unsubscribe(); source.destroy(); target.destroy(); }); }, ); // ═══════════════════════════════════════════════════════════════ // Case 8: AsyncPullable → sync-PollingProxy (Error — sync-poller cannot await) // ═══════════════════════════════════════════════════════════════ describe.each([ ...dataProviderForAsyncPullableToSyncPollingProxyThrows(), ] as Array<[]>)( 'linkTransfers throws error for AsyncPullable to sync PollingProxy test', () => { it('', () => { const mockFlow = { read: jest.fn(async () => 1) }; const source = new AsyncReadTransfer({ flow: mockFlow }); const target = new PollingProxyTransfer({ interval: 100, activated: false }); expect(source.isAsyncPullable).toBe(true); expect(target.isPollingProxy).toBe(true); expect(target.isAsyncPollingProxy).toBe(false); expect(() => linkTransfers(source, target)).toThrow( 'Cannot link AsyncPullable source to sync PollingProxy' ); source.destroy(); target.destroy(); }); }, ); /** * Data provider for error AsyncPullable → sync-PollingProxy. */ function dataProviderForAsyncPullableToSyncPollingProxyThrows(): Array { return [[]]; } // ═══════════════════════════════════════════════════════════════ // Case 9: AsyncPullable → Pushable (Error — needs Bridge/Triggerable) // ═══════════════════════════════════════════════════════════════ describe.each([ ...dataProviderForAsyncPullableToPushableThrows(), ] as Array<[]>)( 'linkTransfers throws error for AsyncPullable to Pushable test', () => { it('', () => { const mockFlow = { read: jest.fn(async () => 1) }; const source = new AsyncReadTransfer({ flow: mockFlow }); const target = new SinkTransfer({ callback: jest.fn() }); expect(source.isAsyncPullable).toBe(true); expect(target.isPushable).toBe(true); expect(() => linkTransfers(source, target)).toThrow( 'Cannot directly link Pullable/AsyncPullable source to Pushable/AsyncPushable target' ); source.destroy(); target.destroy(); }); }, ); /** * Data provider for error AsyncPullable → Pushable. */ function dataProviderForAsyncPullableToPushableThrows(): Array { return [[]]; } describe.each([ ...dataProviderForAsyncPullableToAsyncPushableThrows(), ] as Array<[]>)( 'linkTransfers throws error for AsyncPullable to AsyncPushable test', () => { it('', () => { const mockFlow = { read: jest.fn(async () => 1) }; const source = new AsyncReadTransfer({ flow: mockFlow }); const target = new AsyncSinkTransfer({ callback: jest.fn() as any }); expect(source.isAsyncPullable).toBe(true); expect(target.isAsyncPushable).toBe(true); expect(() => linkTransfers(source, target)).toThrow( 'Cannot directly link Pullable/AsyncPullable source to Pushable/AsyncPushable target' ); source.destroy(); target.destroy(); }); }, ); /** * Data provider for error AsyncPullable → AsyncPushable. */ function dataProviderForAsyncPullableToAsyncPushableThrows(): Array { return [[]]; }