Skip to content
Draft
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
132 changes: 124 additions & 8 deletions src/client/common/pipes/namedPipes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import * as fs from 'fs-extra';
import * as net from 'net';
import * as os from 'os';
import * as path from 'path';
import { PassThrough } from 'stream';
import * as rpc from 'vscode-jsonrpc/node';
import { CancellationError, CancellationToken, Disposable } from 'vscode';
import { traceVerbose } from '../../logging';
Expand All @@ -15,6 +16,9 @@ import { createDeferred } from '../utils/async';
import { noop } from '../utils/misc';

const { XDG_RUNTIME_DIR } = process.env;
const FIFO_READ_RETRY_DELAY_MS = 10;
const FIFO_READ_BUFFER_SIZE = 64 * 1024;

export function generateRandomPipeName(prefix: string): string {
// length of 10 picked because of the name length restriction for sockets
const randomSuffix = crypto.randomBytes(10).toString('hex');
Expand Down Expand Up @@ -149,6 +153,125 @@ class CombinedReader implements rpc.MessageReader {
}
}

// A net.Socket does not surface EOF for a nonblocking FIFO descriptor, so read
// the descriptor directly and close after EOF or owner-signaled completion.
class FifoMessageReader extends rpc.AbstractMessageReader {
private readonly stream = new PassThrough();

private readonly messageReader = new rpc.StreamMessageReader(this.stream, 'utf-8');

private readonly buffer = Buffer.allocUnsafe(FIFO_READ_BUFFER_SIZE);

private readonly disposables: rpc.Disposable[] = [];

private hasReadData = false;

private hasObservedWriter = false;

private disposed = false;

private closePending = false;

constructor(
private readonly pipeName: string,
private readonly fd: number,
private readonly token?: CancellationToken,
) {
super();
this.disposables.push(
this.messageReader.onError((error) => this.fireError(error)),
this.messageReader.onPartialMessage((info) => this.firePartialMessage(info)),
);
}

listen(callback: rpc.DataCallback): rpc.Disposable {
const listener = this.messageReader.listen(callback);
void this.read();
return listener;
}

private async read(): Promise<void> {
while (!this.disposed && !this.closePending) {
try {
const { bytesRead } = await fs.read(this.fd, this.buffer, 0, this.buffer.length, null);
if (bytesRead > 0) {
this.hasReadData = true;
if (!this.stream.write(Buffer.from(this.buffer.subarray(0, bytesRead)))) {
await this.waitForDrain();
}
continue;
}

if (this.hasReadData || this.hasObservedWriter || this.token?.isCancellationRequested) {
this.complete();
return;
}
} catch (error) {
if (this.isRetryableReadError(error)) {
this.hasObservedWriter = true;
} else {
if (!this.disposed) {
this.fireError(error);
this.complete();
}
return;
}
}

await new Promise<void>((resolve) => {
setTimeout(resolve, FIFO_READ_RETRY_DELAY_MS);
});
}
}

private waitForDrain(): Promise<void> {
return new Promise((resolve) => {
const finish = () => {
this.stream.off('drain', finish);
this.stream.off('close', finish);
resolve();
};
this.stream.once('drain', finish);
this.stream.once('close', finish);
});
}

private isRetryableReadError(error: unknown): boolean {
if (!(error instanceof Error) || !('code' in error)) {
return false;
}
const { code } = error as NodeJS.ErrnoException;
return code === 'EAGAIN' || code === 'EWOULDBLOCK';
}

private complete(): void {
if (this.disposed || this.closePending) {
return;
}
this.closePending = true;
this.stream.end();
setImmediate(() => {
if (!this.disposed) {
this.fireClose();
}
});
}

dispose(): void {
if (this.disposed) {
return;
}
this.disposed = true;
this.stream.destroy();
this.messageReader.dispose();
this.disposables.forEach((disposable) => disposable.dispose());
this.disposables.length = 0;
fs.close(this.fd).catch(noop);
fs.unlink(this.pipeName).catch(noop);
super.dispose();
}
}

export async function createReaderPipe(pipeName: string, token?: CancellationToken): Promise<rpc.MessageReader> {
if (isWindows()) {
// windows implementation of FIFO using named pipes
Expand Down Expand Up @@ -189,12 +312,5 @@ export async function createReaderPipe(pipeName: string, token?: CancellationTok
// Intentionally ignored
}
const fd = await fs.open(pipeName, fs.constants.O_RDONLY | fs.constants.O_NONBLOCK);
const socket = new net.Socket({ fd });
const reader = new rpc.SocketMessageReader(socket, 'utf-8');
socket.on('close', () => {
fs.close(fd).catch(noop);
reader.dispose();
});

return reader;
return new FifoMessageReader(pipeName, fd, token);
}
7 changes: 5 additions & 2 deletions src/client/testing/testController/common/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ export async function startRunResultNamedPipe(
dataReceivedCallback: (payload: ExecutionTestPayload) => void,
deferredTillServerClose: Deferred<void>,
cancellationToken?: CancellationToken,
): Promise<string> {
): Promise<{ name: string; dispose: () => void }> {
traceVerbose('Starting Test Result named pipe');
const pipeName: string = generateRandomPipeName('python-test-results');

Expand Down Expand Up @@ -147,7 +147,10 @@ export async function startRunResultNamedPipe(
}),
);

return pipeName;
return {
name: pipeName,
dispose: () => disposable.dispose(),
};
}

interface DiscoveryResultMessage extends Message {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -218,7 +218,8 @@ export class PytestTestDiscoveryAdapter implements ITestDiscoveryAdapter {
} finally {
// Dispose all cancellation handlers and event subscriptions
disposables.forEach((d) => d.dispose());
// Dispose the discovery pipe cancellation token
// Stop the discovery pipe even when the subprocess never opened it.
discoveryPipeCancellation.cancel();
discoveryPipeCancellation.dispose();
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ export class PytestTestExecutionAdapter implements ITestExecutionAdapter {
const cSource = new CancellationTokenSource();
runInstance.token.onCancellationRequested(() => cSource.cancel());

const name = await utils.startRunResultNamedPipe(
const resultPipe = await utils.startRunResultNamedPipe(
dataReceivedCallback, // callback to handle data received
deferredTillServerClose, // deferred to resolve when server closes
cSource.token, // token to cancel
Expand All @@ -66,7 +66,7 @@ export class PytestTestExecutionAdapter implements ITestExecutionAdapter {
await this.runTestsNew(
uri,
testIds,
name,
resultPipe.name,
cSource,
runInstance,
profileKind,
Expand All @@ -76,7 +76,10 @@ export class PytestTestExecutionAdapter implements ITestExecutionAdapter {
project,
);
} finally {
cSource.cancel();
await utils.awaitDeferredWithTimeout(deferredTillServerClose, utils.RESULT_PIPE_DRAIN_TIMEOUT_MS);
resultPipe.dispose();
cSource.dispose();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -204,7 +204,8 @@ export class UnittestTestDiscoveryAdapter implements ITestDiscoveryAdapter {
traceVerbose(`Cleaning up unittest discovery resources for workspace ${uri.fsPath}`);
// Dispose all cancellation handlers and event subscriptions
disposables.forEach((d) => d.dispose());
// Dispose the discovery pipe cancellation token
// Stop the discovery pipe even when the subprocess never opened it.
discoveryPipeCancellation.cancel();
discoveryPipeCancellation.dispose();
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ export class UnittestTestExecutionAdapter implements ITestExecutionAdapter {
};
const cSource = new CancellationTokenSource();
runInstance.token.onCancellationRequested(() => cSource.cancel());
const name = await utils.startRunResultNamedPipe(
const resultPipe = await utils.startRunResultNamedPipe(
dataReceivedCallback, // callback to handle data received
deferredTillServerClose, // deferred to resolve when server closes
cSource.token, // token to cancel
Expand All @@ -79,7 +79,7 @@ export class UnittestTestExecutionAdapter implements ITestExecutionAdapter {
await this.runTestsNew(
uri,
testIds,
name,
resultPipe.name,
cSource,
runInstance,
profileKind,
Expand All @@ -90,7 +90,10 @@ export class UnittestTestExecutionAdapter implements ITestExecutionAdapter {
} catch (error) {
traceError(`Error in running unittest tests: ${error}`);
} finally {
cSource.cancel();
await utils.awaitDeferredWithTimeout(deferredTillServerClose, utils.RESULT_PIPE_DRAIN_TIMEOUT_MS);
resultPipe.dispose();
cSource.dispose();
}
}

Expand Down
Loading
Loading