Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,17 @@ public class InputStreamMonitor {
*/
private volatile boolean fClosed = false;

/**
* Whether {@link #closeInputStream()} was called. Guarded by {@link #fLock}.
*/
private boolean fCloseRequested;

/**
* Queued after the last data to make the writer thread close the stream once
* all previously queued data is written. Compared by identity.
*/
private static final byte[] CLOSE_MARKER = new byte[0];

/**
* The charset of the input stream.
*/
Expand Down Expand Up @@ -109,10 +120,10 @@ public void write(byte[] data, int offset, int length) {
}

private void write(byte[] copy) {
if (fClosed) {
return; // drop data;
}
synchronized (fLock) {
if (fClosed || fCloseRequested) {
return; // drop data;
}
fQueue.offer(copy);
fLock.notifyAll();
}
Expand Down Expand Up @@ -158,7 +169,7 @@ public void close() {
private void write() {
try {
try {
while (fThread != null) {
while (fThread != null && !fClosed) {
writeNext();
}
} finally {
Expand All @@ -179,6 +190,11 @@ private void write() {
private void writeNext() throws IOException {
while (!fQueue.isEmpty() && !fClosed) {
byte[] data = fQueue.poll();
if (data == CLOSE_MARKER) {
fClosed = true;
fStream.close();
return;
}
fStream.write(data);
fStream.flush();
}
Expand All @@ -187,7 +203,7 @@ private void writeNext() throws IOException {
// Queue could receive more input between last empty check and
// lock acquire. See https://bugs.eclipse.org/550834
// Use while instead of if to guard against spurious wakeups.
while (fQueue.isEmpty()) {
while (fQueue.isEmpty() && !fClosed) {
fLock.wait();
}
}
Expand All @@ -197,19 +213,31 @@ private void writeNext() throws IOException {

/**
* Closes the output stream attached to the standard input stream of this
* monitor's process.
* monitor's process. Data queued before this call is written first and the
* stream is closed afterwards by the writer thread, also if monitoring is
* started only later. Without pending data and a writer thread the stream is
* closed immediately.
*
* @exception IOException if an exception occurs closing the input stream or
* stream is already closed
*/
public void closeInputStream() throws IOException {
if (!fClosed) {
synchronized (fLock) {
if (fClosed || fCloseRequested) {
throw new IOException();
}
fCloseRequested = true;
if (fThread != null || !fQueue.isEmpty()) {
fQueue.offer(CLOSE_MARKER);
fLock.notifyAll();
return;
}
// Set under the lock so a writer thread started concurrently by
// startMonitoring() sees it and does not wait for data forever.
fClosed = true;
fStream.close();
} else {
throw new IOException();
fLock.notifyAll();
}

fStream.close();
}
}

Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,15 @@

import static org.assertj.core.api.Assertions.assertThat;

import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.OutputStream;
import java.io.PipedInputStream;
import java.io.PipedOutputStream;
import java.nio.charset.StandardCharsets;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Supplier;

import org.eclipse.core.runtime.ILog;
Expand Down Expand Up @@ -82,6 +86,89 @@ private void waitForElementsInStream(PipedInputStream sysin, int numberOfElement
assertThat(sysin.available()).isEqualTo(numberOfElements);
}

/**
* Data written before {@link InputStreamMonitor#closeInputStream()} must reach
* the stream before it is closed.
*/
@Test
@SuppressWarnings("resource")
public void testCloseInputStreamWritesPendingData() throws Exception {
ByteArrayOutputStream written = new ByteArrayOutputStream();
AtomicBoolean writtenAfterClose = new AtomicBoolean();
AtomicBoolean closed = new AtomicBoolean();
OutputStream slowStream = new OutputStream() {
@Override
public void write(int b) throws IOException {
write(new byte[] { (byte) b }, 0, 1);
}

@Override
public void write(byte[] b, int off, int len) throws IOException {
if (closed.get()) {
writtenAfterClose.set(true);
}
try {
Thread.sleep(5);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
written.write(b, off, len);
}

@Override
public void close() {
closed.set(true);
}
};
InputStreamMonitor monitor = new InputStreamMonitor(slowStream);
try {
monitor.startMonitoring();
byte[] chunk = new byte[] { 1, 2, 3, 4 };
int chunks = 20;
for (int i = 0; i < chunks; i++) {
monitor.write(chunk, 0, chunk.length);
}
monitor.closeInputStream();
TestUtil.waitWhile(() -> !closed.get(), CONDITION_TIMEOUT_IN_MILLIS);
assertThat(closed.get()).withFailMessage("stream not closed").isTrue();
assertThat(written.size()).as("bytes written before close").isEqualTo(chunks * chunk.length);
assertThat(writtenAfterClose.get()).as("data written after close").isFalse();
} finally {
monitor.close();
}
}

/**
* Data written before monitoring is started and followed by
* {@link InputStreamMonitor#closeInputStream()} must still reach the stream
* before it is closed.
*/
@Test
public void testCloseInputStreamBeforeStartWritesPendingData() throws Exception {
AtomicInteger numClosed = new AtomicInteger();
AtomicInteger bytesWrittenAtClose = new AtomicInteger(-1);
ByteArrayOutputStream written = new ByteArrayOutputStream() {
@Override
public void close() {
bytesWrittenAtClose.set(size());
numClosed.incrementAndGet();
}
};
byte[] content = "0 8 1\n2 7 1\n".getBytes(StandardCharsets.US_ASCII);
InputStreamMonitor monitor = new InputStreamMonitor(written);
try {
monitor.write(content, 0, content.length);
monitor.closeInputStream();
monitor.startMonitoring();
TestUtil.waitWhile(() -> numClosed.get() == 0, CONDITION_TIMEOUT_IN_MILLIS);
assertThat(numClosed.get()).as("stream close count").isEqualTo(1);
assertThat(bytesWrittenAtClose.get()).as("bytes written before close").isEqualTo(content.length);
assertThat(written.toByteArray()).as("written content").isEqualTo(content);
} finally {
monitor.close();
}
}

/**
* Test that passing <code>null</code> as charset does not raise exceptions.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,13 @@ public class MockProcess extends Process {
public static final int RUN_FOREVER = -1;

/** Mockup processe's standard streams. */
private final ByteArrayOutputStream stdin = new ByteArrayOutputStream();
private final ByteArrayOutputStream stdin = new ByteArrayOutputStream() {
@Override
public void close() {
stdinClosed = true;
}
};
private volatile boolean stdinClosed;
private final InputStream stdout;
private final InputStream stderr;

Expand Down Expand Up @@ -215,6 +221,14 @@ public synchronized byte[] getReceivedInput() {
return content;
}

/**
* @return whether the standard input stream of this process was closed, i.e.
* the process would read EOF
*/
public boolean isStdinClosed() {
return stdinClosed;
}

@Override
public OutputStream getOutputStream() {
return stdin;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -586,6 +586,39 @@ public void testBinaryInputFromFile() throws Exception {
assertThat(receivedInput).as("received input").isEqualTo(input);
}

/**
* Test that a process reading its standard input from a file receives EOF
* after the whole file content.
*/
@Test
Comment thread
anydoby marked this conversation as resolved.
public void testInputFromFileSendsEof() throws Exception {
byte[] input = "0 8 1\n2 7 1\n3 4 2\n".getBytes(StandardCharsets.US_ASCII);
String consoleEncoding = StandardCharsets.UTF_8.name();

final File inFile = createTmpFile("testinput-eof.txt");
Files.write(inFile.toPath(), input);
final MockProcess mockProcess = new MockProcess(MockProcess.RUN_FOREVER);
try {
Map<String, Object> launchConfigAttributes = new HashMap<>();
launchConfigAttributes.put(DebugPlugin.ATTR_CONSOLE_ENCODING, consoleEncoding);
launchConfigAttributes.put(IDebugUIConstants.ATTR_CAPTURE_STDIN_FILE, inFile.getCanonicalPath());
launchConfigAttributes.put(IDebugUIConstants.ATTR_CAPTURE_IN_CONSOLE, false);
final IProcess process = mockProcess.toRuntimeProcess("inputFromFileEof", launchConfigAttributes);
final ProcessConsole console = new ProcessConsole(process, new ConsoleColorProvider(), consoleEncoding);
try {
console.initialize();
waitWhile(() -> !mockProcess.isStdinClosed(), () -> "Process stdin was not closed after input file was read.");
} finally {
console.destroy();
}
} finally {
mockProcess.destroy();
}

byte[] receivedInput = mockProcess.getReceivedInput();
assertThat(receivedInput).as("received input").isEqualTo(input);
}

/**
* Test that console name updates (elapsed time) only happen for visible
* consoles. Hidden consoles should not update their name. When a hidden
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@
import org.eclipse.debug.core.model.IProcess;
import org.eclipse.debug.core.model.IStreamMonitor;
import org.eclipse.debug.core.model.IStreamsProxy;
import org.eclipse.debug.core.model.IStreamsProxy2;
Comment thread
anydoby marked this conversation as resolved.
import org.eclipse.debug.core.sourcelookup.containers.LocalFileStorage;
import org.eclipse.debug.internal.core.IInternalDebugCoreConstants;
import org.eclipse.debug.internal.ui.DebugPluginImages;
Expand Down Expand Up @@ -1022,6 +1023,7 @@ protected IStatus run(IProgressMonitor monitor) {
if (readingStream == null || isStreamClosed()) {
return monitor.isCanceled() ? Status.CANCEL_STATUS : Status.OK_STATUS;
}
boolean endOfInput = false;
if (streamsProxy instanceof IBinaryStreamsProxy proxy) {
// Pass data without processing. The preferred variant. There is no need for
// this job to know about encodings.
Expand All @@ -1037,6 +1039,7 @@ protected IStatus run(IProgressMonitor monitor) {
proxy.write(buffer, 0, bytesRead);
}
}
endOfInput = bytesRead < 0;
} catch (IOException e) {
if (!isStreamClosed()) {
DebugUIPlugin.log(e);
Expand All @@ -1062,12 +1065,23 @@ protected IStatus run(IProgressMonitor monitor) {
streamsProxy.write(s);
}
}
endOfInput = charRead < 0;
} catch (IOException e) {
if (!isStreamClosed()) {
DebugUIPlugin.log(e);
}
}
}
// Input redirected from a file must reach the process as EOF once the file is
// exhausted, otherwise programs reading stdin until EOF never terminate.
if (endOfInput && readingStream != fUserInput && !isStreamClosed()
&& streamsProxy instanceof IStreamsProxy2 proxy2) {
try {
proxy2.closeInputStream();
} catch (IOException e) {
// already closed, e.g. process terminated
}
}
return monitor.isCanceled() ? Status.CANCEL_STATUS : Status.OK_STATUS;
}
}
Expand Down
Loading