Skip to content

Commit cca7db5

Browse files
committed
Preserve Windows terminal output through teardown
1 parent 872bf9a commit cca7db5

5 files changed

Lines changed: 124 additions & 37 deletions

File tree

‎CMakeLists.txt‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -674,13 +674,23 @@ if(WIN32)
674674
${CMAKE_CURRENT_SOURCE_DIR}/test/integration_tests/*Test.cpp
675675
${CMAKE_CURRENT_SOURCE_DIR}/test/unit_tests/*Test.cpp)
676676

677+
add_executable(
678+
et-conpty-output-helper
679+
test/helpers/ConPtyOutputHelper.cpp)
680+
677681
# Same test set as Unix (see the else(WIN32) branch below): every test is
678682
# expected to compile and pass on Windows. Platform-impossible cases
679683
# SKIP at runtime with a reason instead of being excluded here.
680684
add_executable(et-test ${WINDOWS_TEST_SRCS} test/Main.cpp)
681685
target_include_directories(
682686
et-test PRIVATE test test/integration_tests test/unit_tests)
683-
add_dependencies(et-test generated-code TerminalCommon et-lib HtmCommon)
687+
add_dependencies(
688+
et-test
689+
generated-code
690+
TerminalCommon
691+
et-lib
692+
HtmCommon
693+
et-conpty-output-helper)
684694
target_link_libraries(
685695
et-test
686696
PRIVATE

‎src/terminal/PseudoUserTerminalWindows.hpp‎

Lines changed: 17 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -221,7 +221,6 @@ class PseudoUserTerminal : public UserTerminal {
221221
char bytes[16 * 1024];
222222
DWORD count = 0;
223223
HANDLE output = static_cast<HANDLE>(outputRead);
224-
constexpr size_t maxPendingOutput = 256 * 1024;
225224
// Exits only when ReadFile fails, i.e. once ClosePseudoConsole has
226225
// released the output pipe (see cleanup()). The pending-output cap
227226
// below is bypassed once closing is set, so a stalled consumer can
@@ -230,21 +229,20 @@ class PseudoUserTerminal : public UserTerminal {
230229
while (true) {
231230
{
232231
std::unique_lock<std::mutex> guard(pendingMutex);
233-
pendingDrained.wait(guard, [this, maxPendingOutput]() {
234-
return closing || pendingOutput.size() < maxPendingOutput;
232+
pendingDrained.wait(guard, [this]() {
233+
return closing || pendingOutput.size() < MAX_PENDING_OUTPUT;
235234
});
236235
}
237236
if (!ReadFile(output, bytes, sizeof(bytes), &count, NULL) ||
238237
count == 0) {
239238
break;
240239
}
241240
std::lock_guard<std::mutex> guard(pendingMutex);
242-
if (closing) {
243-
// Nobody will call drainOutput() again; keep reading so the host's
244-
// writes never block on a full pipe, but drop earlier bytes so
245-
// this can't grow unbounded while we wait for it to exit.
246-
pendingOutput.clear();
247-
}
241+
// cleanup() terminates the process tree before closing the
242+
// pseudoconsole, so bypassing the normal cap during teardown has a
243+
// bounded producer. Preserve everything already queued while this
244+
// thread drains the output pipe; clearing it here loses terminal
245+
// output that the router may still deliver after reconnect.
248246
pendingOutput.append(bytes, count);
249247
}
250248
});
@@ -256,15 +254,9 @@ class PseudoUserTerminal : public UserTerminal {
256254
virtual void runTerminal() override {}
257255

258256
virtual void cleanup() override {
259-
// Stop enforcing the pending-output cap and wake the drain thread in
260-
// case it is currently parked there. Without this, a stalled consumer
261-
// (router backpressure, or nobody left to call drainOutput()) could
262-
// leave the ConPTY output pipe undrained while the host process below
263-
// still has queued output to write; ClosePseudoConsole then blocks
264-
// waiting for a host that is itself blocked on a full pipe, hanging
265-
// forever.
266-
closing = true;
267-
pendingDrained.notify_all();
257+
// Interrupt input and the process tree before lifting the output cap. This
258+
// bounds how much final output teardown can append while preserving every
259+
// byte that was already queued for the router.
268260
if (inputWrite != INVALID_HANDLE_VALUE) {
269261
CloseHandle(static_cast<HANDLE>(inputWrite));
270262
inputWrite = INVALID_HANDLE_VALUE;
@@ -284,6 +276,12 @@ class PseudoUserTerminal : public UserTerminal {
284276
CloseHandle(static_cast<HANDLE>(jobHandle));
285277
jobHandle = nullptr;
286278
}
279+
// Wake the reader if router backpressure parked it at the cap. It must
280+
// continue draining while ClosePseudoConsole runs: older Windows waits
281+
// for the pseudoconsole host, which can itself be blocked writing final
282+
// output to a full pipe.
283+
closing = true;
284+
pendingDrained.notify_all();
287285
if (hPC != nullptr) {
288286
// Blocks until the ConPTY host process exits. The drain thread above
289287
// keeps calling ReadFile the whole time (closing bypasses its cap), so
@@ -368,6 +366,7 @@ class PseudoUserTerminal : public UserTerminal {
368366
}
369367

370368
protected:
369+
static constexpr size_t MAX_PENDING_OUTPUT = 256 * 1024;
371370
void* hPC;
372371
void* inputWrite;
373372
void* outputRead;
Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
#include <windows.h>
2+
3+
#include <array>
4+
#include <cstddef>
5+
6+
int main() {
7+
// Open the attached console explicitly. Test runners may give a child
8+
// inherited standard handles of their own, while CONOUT$ always targets the
9+
// pseudoconsole selected through PROC_THREAD_ATTRIBUTE_PSEUDOCONSOLE.
10+
HANDLE output =
11+
CreateFileW(L"CONOUT$", GENERIC_WRITE, FILE_SHARE_READ | FILE_SHARE_WRITE,
12+
nullptr, OPEN_EXISTING, 0, nullptr);
13+
if (output == nullptr || output == INVALID_HANDLE_VALUE) {
14+
return 1;
15+
}
16+
17+
std::array<char, 16 * 1024> chunk;
18+
chunk.fill('A');
19+
constexpr size_t totalBytes = 1024 * 1024;
20+
size_t totalWritten = 0;
21+
while (totalWritten < totalBytes) {
22+
size_t chunkOffset = 0;
23+
while (chunkOffset < chunk.size()) {
24+
DWORD written = 0;
25+
if (!WriteFile(output, chunk.data() + chunkOffset,
26+
static_cast<DWORD>(chunk.size() - chunkOffset), &written,
27+
nullptr) ||
28+
written == 0) {
29+
CloseHandle(output);
30+
return 2;
31+
}
32+
chunkOffset += written;
33+
totalWritten += written;
34+
}
35+
}
36+
CloseHandle(output);
37+
return 0;
38+
}

‎test/integration_tests/TerminalTest.cpp‎

Lines changed: 19 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -706,6 +706,12 @@ class PartialFailurePipeSocketHandler : public PipeSocketHandler {
706706
outputToFail = output;
707707
partialFd = -1;
708708
partialWriteDone = false;
709+
targetWriteFailed = false;
710+
}
711+
712+
bool didFailTargetWrite() {
713+
lock_guard<mutex> guard(failureMutex);
714+
return targetWriteFailed;
709715
}
710716

711717
ssize_t write(int fd, const void* buf, size_t count) override {
@@ -724,6 +730,8 @@ class PartialFailurePipeSocketHandler : public PipeSocketHandler {
724730
sendPartial = true;
725731
} else if (partialWriteDone && fd == partialFd) {
726732
outputToFail.clear();
733+
partialFd = -1;
734+
targetWriteFailed = true;
727735
failWrite = true;
728736
}
729737
}
@@ -744,6 +752,7 @@ class PartialFailurePipeSocketHandler : public PipeSocketHandler {
744752
string outputToFail;
745753
int partialFd = -1;
746754
bool partialWriteDone = false;
755+
bool targetWriteFailed = false;
747756
};
748757
#endif
749758

@@ -906,14 +915,16 @@ TEST_CASE_METHOD(EndToEndTestFixture, "WindowsOutputSurvivesFailedWrite",
906915
std::chrono::steady_clock::now() < terminalDeadline) {
907916
std::this_thread::sleep_for(std::chrono::milliseconds(10));
908917
}
918+
const bool terminalWasSetup = fakeUserTerminal->isSetup();
909919

910920
const string marker = "ET_WINDOWS_PENDING_OUTPUT";
911-
if (fakeUserTerminal->isSetup()) {
921+
if (terminalWasSetup) {
912922
routerSocketHandler->failNextWriteContaining(marker);
913923
fakeUserTerminal->simulateTerminalResponse(marker);
914924
}
915925

916926
string received;
927+
bool terminalSurvivedReconnect = false;
917928
const auto outputDeadline =
918929
std::chrono::steady_clock::now() + std::chrono::seconds(30);
919930
while (fakeUserTerminal->isSetup() &&
@@ -922,6 +933,8 @@ TEST_CASE_METHOD(EndToEndTestFixture, "WindowsOutputSurvivesFailedWrite",
922933
received += fakeConsole->getTerminalData(1);
923934
}
924935
if (received == marker) {
936+
terminalSurvivedReconnect =
937+
fakeUserTerminal->isSetup() && !fakeUserTerminal->wasCleanedUp();
925938
break;
926939
}
927940
std::this_thread::sleep_for(std::chrono::milliseconds(10));
@@ -935,8 +948,12 @@ TEST_CASE_METHOD(EndToEndTestFixture, "WindowsOutputSurvivesFailedWrite",
935948
handlerThread.join();
936949
handler.reset();
937950

938-
REQUIRE(fakeUserTerminal->isSetup());
951+
REQUIRE(terminalWasSetup);
952+
REQUIRE(routerSocketHandler->didFailTargetWrite());
939953
REQUIRE(received == marker);
954+
REQUIRE(terminalSurvivedReconnect);
955+
REQUIRE(fakeUserTerminal->wasCleanedUp());
956+
REQUIRE_FALSE(fakeUserTerminal->isSetup());
940957
}
941958
#endif
942959

‎test/unit_tests/PseudoUserTerminalWindowsTest.cpp‎

Lines changed: 39 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
#include <windows.h>
44

5+
#include <array>
56
#include <chrono>
67
#include <future>
78
#include <thread>
@@ -13,10 +14,12 @@ namespace et {
1314

1415
class OutputPressureTerminal : public PseudoUserTerminal {
1516
public:
16-
size_t queuedOutputSize() {
17+
std::string queuedOutput() {
1718
std::lock_guard<std::mutex> guard(pendingMutex);
18-
return pendingOutput.size();
19+
return pendingOutput;
1920
}
21+
22+
static constexpr size_t outputLimit() { return MAX_PENDING_OUTPUT; }
2023
};
2124

2225
TEST_CASE("PseudoUserTerminal cleanup falls back without a job",
@@ -53,35 +56,47 @@ TEST_CASE("PseudoUserTerminal cleanup falls back without a job",
5356
// parked waiting for drainOutput() to free up its pending-output cap, and
5457
// nobody calls drainOutput() again during shutdown), the ConPTY host can
5558
// block writing to a full pipe and ClosePseudoConsole waits for that host
56-
// to exit forever. Heap-allocate the terminal and leak it (rather than
57-
// join/delete) if cleanup() does not return in time, since a hung cleanup()
58-
// thread could otherwise touch freed memory.
59+
// to exit forever. Teardown must continue draining without discarding output
60+
// that was already queued for the router. Heap-allocate the terminal and leak
61+
// it (rather than join/delete) if cleanup() does not return in time, since a
62+
// hung cleanup() thread could otherwise touch freed memory.
5963
TEST_CASE(
6064
"PseudoUserTerminal cleanup does not hang under queued, undrained output",
6165
"[PseudoUserTerminal][windows]") {
6266
const char* previousShell = ::getenv("SHELL");
6367
const bool hadShell = previousShell != nullptr;
6468
const std::string savedShell = hadShell ? previousShell : "";
65-
REQUIRE(_putenv_s("SHELL", "cmd.exe") == 0);
69+
std::array<wchar_t, 32768> testExecutable;
70+
const DWORD testExecutableLength =
71+
GetModuleFileNameW(nullptr, testExecutable.data(),
72+
static_cast<DWORD>(testExecutable.size()));
73+
REQUIRE(testExecutableLength > 0);
74+
REQUIRE(testExecutableLength < testExecutable.size());
75+
const fs::path helperPath =
76+
fs::path(std::wstring(testExecutable.data(), testExecutableLength))
77+
.parent_path() /
78+
"et-conpty-output-helper.exe";
79+
REQUIRE(fs::is_regular_file(helperPath));
80+
REQUIRE(fs::is_directory(helperPath.parent_path()));
81+
REQUIRE(_putenv_s("SHELL", helperPath.string().c_str()) == 0);
6682
auto* term = new OutputPressureTerminal();
6783
const int setupResult = term->setup(-1);
6884
REQUIRE(_putenv_s("SHELL", savedShell.c_str()) == 0);
6985
REQUIRE(setupResult == 0);
7086

71-
// Flood stdout continuously and never call drainOutput(), modeling a
72-
// stalled router consumer. This drives the drain thread past its
73-
// pending-output cap and fills the underlying pipe before teardown.
74-
term->writeInput(
75-
"for /L %i in (1,1,1000000) do @echo "
76-
"AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"
77-
"\r\n");
87+
// The helper writes a finite 1 MiB directly to its ConPTY stdout. Never
88+
// call drainOutput() before cleanup, modeling a stalled router consumer
89+
// without relying on shell startup, command parsing, or loop throughput.
7890
const auto pressureDeadline =
79-
std::chrono::steady_clock::now() + std::chrono::seconds(5);
80-
while (term->queuedOutputSize() < 256 * 1024 &&
91+
std::chrono::steady_clock::now() + std::chrono::seconds(15);
92+
string queuedBeforeCleanup;
93+
while (queuedBeforeCleanup.size() < OutputPressureTerminal::outputLimit() &&
8194
std::chrono::steady_clock::now() < pressureDeadline) {
8295
std::this_thread::sleep_for(std::chrono::milliseconds(20));
96+
queuedBeforeCleanup = term->queuedOutput();
8397
}
84-
const bool capReached = term->queuedOutputSize() >= 256 * 1024;
98+
const bool capReached =
99+
queuedBeforeCleanup.size() >= OutputPressureTerminal::outputLimit();
85100

86101
auto done = std::make_shared<std::promise<void>>();
87102
std::future<void> doneFuture = done->get_future();
@@ -91,14 +106,22 @@ TEST_CASE(
91106
});
92107

93108
const auto status = doneFuture.wait_for(std::chrono::seconds(15));
109+
bool queuedOutputPreserved = false;
94110
if (status == std::future_status::ready) {
95111
cleanupThread.join();
112+
const string queuedAfterCleanup = term->drainOutput();
113+
queuedOutputPreserved =
114+
queuedAfterCleanup.size() >= queuedBeforeCleanup.size() &&
115+
std::equal(queuedBeforeCleanup.begin(), queuedBeforeCleanup.end(),
116+
queuedAfterCleanup.begin());
96117
delete term;
97118
} else {
98119
cleanupThread.detach();
99120
}
121+
INFO("queued before cleanup: " << queuedBeforeCleanup.size());
100122
REQUIRE(capReached);
101123
REQUIRE(status == std::future_status::ready);
124+
REQUIRE(queuedOutputPreserved);
102125
}
103126

104127
} // namespace et

0 commit comments

Comments
 (0)