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
38 changes: 26 additions & 12 deletions src/data/operation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -101,19 +101,24 @@ export function runOperation<T>(
let settled = false;
let cleanup: Cleanup;
const unsubscribers: Array<() => void> = [];
let finishStart!: () => void;
const startFinished = new Promise<void>((resolve) => { finishStart = resolve; });

const finish = (error: TradingViewError | undefined, value?: T) => {
if (settled) return;
settled = true;
clearTimeout(timer);
signal?.removeEventListener('abort', onAbort);
for (const off of unsubscribers) off();
try {
cleanup?.();
} catch {
// Cleanup must never hide the result.
}
connection.release().catch(() => {}).then(() => {
// `start` may resolve synchronously, before it has returned its cleanup.
startFinished.then(() => {
try {
cleanup?.();
} catch {
// Cleanup must never hide the result.
}
return connection.release().catch(() => {});
}).then(() => {
if (error) reject(error);
else resolve(value as T);
});
Expand All @@ -138,6 +143,8 @@ export function runOperation<T>(
});
} catch (error) {
finish(toTradingViewError(error, 'INVALID_ARGUMENT'));
} finally {
finishStart();
}
});
}
Expand Down Expand Up @@ -195,6 +202,8 @@ export function startWatcher<W extends Watcher>(
let active = true;
let started = false;
let cleanup: Cleanup;
let finishStart!: () => void;
const startFinished = new Promise<void>((resolve) => { finishStart = resolve; });
let stopPromise: Promise<void> | undefined;
let resolveClosed!: () => void;
const closed = new Promise<void>((resolve) => { resolveClosed = resolve; });
Expand All @@ -220,12 +229,15 @@ export function startWatcher<W extends Watcher>(
clearTimeout(timer);
signal?.removeEventListener('abort', onAbort);
for (const off of unsubscribers) off();
try {
cleanup?.();
} catch {
// Ignore cleanup failures.
}
stopPromise = connection.release().catch(() => {}).then(() => resolveClosed());
// A synchronous ready/fail can stop before `start` returns its cleanup.
stopPromise = startFinished.then(() => {
try {
cleanup?.();
} catch {
// Ignore cleanup failures.
}
return connection.release().catch(() => {});
}).then(() => resolveClosed());
return stopPromise;
};

Expand Down Expand Up @@ -273,6 +285,8 @@ export function startWatcher<W extends Watcher>(
});
} catch (error) {
fail(error);
} finally {
finishStart();
}

return startPromise;
Expand Down
59 changes: 59 additions & 0 deletions tests/unit/operation.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
import { describe, expect, it, vi } from 'vitest';
import { TradingViewClient } from '../../src/client/client.js';
import { runOperation, startWatcher } from '../../src/data/operation.js';
import { FakeServer } from '../helpers/fake-server.js';

function sharedClient() {
const server = new FakeServer();
return new TradingViewClient({ transport: server.transport });
}

describe('synchronous operation lifecycle', () => {
it('runs cleanup before resolving a one-shot operation', async () => {
const client = sharedClient();
const cleanup = vi.fn();
try {
const value = await runOperation({ client }, 'immediate', ({ resolve }) => {
resolve(42);
return cleanup;
});
expect(value).toBe(42);
expect(cleanup).toHaveBeenCalledOnce();
expect(client.isClosed).toBe(false);
} finally {
await client.close();
}
});

it('runs cleanup after a watcher becomes ready synchronously', async () => {
const client = sharedClient();
const cleanup = vi.fn();
try {
const watcher = await startWatcher({ client }, 'immediate', {}, ({ ready }) => {
ready();
return cleanup;
}, (base) => base);
await watcher.stop();
await watcher.closed;
expect(cleanup).toHaveBeenCalledOnce();
expect(client.isClosed).toBe(false);
} finally {
await client.close();
}
});

it('runs cleanup after a watcher fails synchronously', async () => {
const client = sharedClient();
const cleanup = vi.fn();
try {
await expect(startWatcher({ client }, 'immediate', {}, ({ fail }) => {
fail(new Error('startup failed'));
return cleanup;
}, (base) => base)).rejects.toThrow('startup failed');
expect(cleanup).toHaveBeenCalledOnce();
expect(client.isClosed).toBe(false);
} finally {
await client.close();
}
});
});
Loading