diff --git a/debug/org.eclipse.debug.core/core/org/eclipse/debug/internal/core/InputStreamMonitor.java b/debug/org.eclipse.debug.core/core/org/eclipse/debug/internal/core/InputStreamMonitor.java index 19b9c8343ad..776bc765e35 100644 --- a/debug/org.eclipse.debug.core/core/org/eclipse/debug/internal/core/InputStreamMonitor.java +++ b/debug/org.eclipse.debug.core/core/org/eclipse/debug/internal/core/InputStreamMonitor.java @@ -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. */ @@ -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(); } @@ -158,7 +169,7 @@ public void close() { private void write() { try { try { - while (fThread != null) { + while (fThread != null && !fClosed) { writeNext(); } } finally { @@ -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(); } @@ -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(); } } @@ -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(); } } diff --git a/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/InputStreamMonitorTests.java b/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/InputStreamMonitorTests.java index 89dc865df18..5b93e519a04 100644 --- a/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/InputStreamMonitorTests.java +++ b/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/InputStreamMonitorTests.java @@ -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; @@ -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 null as charset does not raise exceptions. */ diff --git a/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/MockProcess.java b/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/MockProcess.java index 14be9f3fab1..e7be3a37af9 100644 --- a/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/MockProcess.java +++ b/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/MockProcess.java @@ -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; @@ -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; diff --git a/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/ProcessConsoleTests.java b/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/ProcessConsoleTests.java index 88c58f845b7..3a6d09bd4e0 100644 --- a/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/ProcessConsoleTests.java +++ b/debug/org.eclipse.debug.tests/src/org/eclipse/debug/tests/console/ProcessConsoleTests.java @@ -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 + 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 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 diff --git a/debug/org.eclipse.debug.ui/ui/org/eclipse/debug/internal/ui/views/console/ProcessConsole.java b/debug/org.eclipse.debug.ui/ui/org/eclipse/debug/internal/ui/views/console/ProcessConsole.java index b577268c14a..7ae8c565346 100644 --- a/debug/org.eclipse.debug.ui/ui/org/eclipse/debug/internal/ui/views/console/ProcessConsole.java +++ b/debug/org.eclipse.debug.ui/ui/org/eclipse/debug/internal/ui/views/console/ProcessConsole.java @@ -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; import org.eclipse.debug.core.sourcelookup.containers.LocalFileStorage; import org.eclipse.debug.internal.core.IInternalDebugCoreConstants; import org.eclipse.debug.internal.ui.DebugPluginImages; @@ -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. @@ -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); @@ -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; } }