diff --git a/src/compiler/propagation/type_primitive_pipe.cc b/src/compiler/propagation/type_primitive_pipe.cc index 55f55ef9b..673212e8e 100644 --- a/src/compiler/propagation/type_primitive_pipe.cc +++ b/src/compiler/propagation/type_primitive_pipe.cc @@ -25,6 +25,7 @@ TYPE_PRIMITIVE_ANY(close) TYPE_PRIMITIVE_ANY(create_pipe) TYPE_PRIMITIVE_ANY(fd_to_pipe) TYPE_PRIMITIVE_ANY(write) +TYPE_PRIMITIVE_ANY(write_result) TYPE_PRIMITIVE_ANY(read) TYPE_PRIMITIVE_ANY(fork) TYPE_PRIMITIVE_ANY(fork2) diff --git a/src/event_sources/event_win.cc b/src/event_sources/event_win.cc index 32490f70e..763e88d09 100644 --- a/src/event_sources/event_win.cc +++ b/src/event_sources/event_win.cc @@ -24,6 +24,16 @@ namespace toit { WindowsEventSource* WindowsEventSource::instance_ = null; +void WindowsOverlapped::cancel_and_wait(HANDLE handle) { + if (!started_) return; + CancelIoEx(handle, &overlapped_); + // Wait even if CancelIoEx didn't find the operation: it may be in the middle + // of completing. Returns immediately if it has completed, and the + // cancellation keeps the wait on the event thread short otherwise. + DWORD count; + GetOverlappedResult(handle, &overlapped_, &count, TRUE); +} + class WindowsEventThread; class WindowsResourceEvent { public: diff --git a/src/event_sources/event_win.h b/src/event_sources/event_win.h index 7f9f1d972..91fcf8dbe 100644 --- a/src/event_sources/event_win.h +++ b/src/event_sources/event_win.h @@ -35,6 +35,37 @@ class WindowsResource : public Resource { virtual bool is_event_enabled(HANDLE event) { return true; } }; +// An OVERLAPPED structure, and whether the last operation issued with it was +// started. +// +// Closing a handle cancels its pending operations, but their completion can +// still be written to the OVERLAPPED structure after the close returns. Call +// cancel_and_wait before closing the handle or the event, and before freeing +// the structure. +class WindowsOverlapped { + public: + OVERLAPPED* get() { return &overlapped_; } + HANDLE event() const { return overlapped_.hEvent; } + void set_event(HANDLE event) { overlapped_.hEvent = event; } + + // Records the result of an overlapped call, like ReadFile or WSARecv. + // Returns whether the operation started, which includes completing + // synchronously. Doesn't change the last error. + bool issued(bool success, DWORD error) { + started_ = success || error == ERROR_IO_PENDING; + return started_; + } + + // Cancels the last operation and waits for it to complete, unless it failed + // to start: the structure isn't meaningful after a synchronous failure, and + // waiting could block forever. + void cancel_and_wait(HANDLE handle); + + private: + OVERLAPPED overlapped_{}; + bool started_ = false; +}; + class WindowsEventThread; class WindowsResourceEvent; diff --git a/src/primitive.h b/src/primitive.h index adcedc05a..1dddff596 100644 --- a/src/primitive.h +++ b/src/primitive.h @@ -792,6 +792,7 @@ namespace toit { PRIMITIVE(fork2, 10) \ PRIMITIVE(fd, 1) \ PRIMITIVE(is_a_tty, 1) \ + PRIMITIVE(write_result, 2) \ #define MODULE_STDIO(PRIMITIVE) \ PRIMITIVE(stdin_init, 0) \ diff --git a/src/resources/pipe_posix.cc b/src/resources/pipe_posix.cc index 50e4ba9d7..4dcf25716 100644 --- a/src/resources/pipe_posix.cc +++ b/src/resources/pipe_posix.cc @@ -243,6 +243,13 @@ PRIMITIVE(write) { return Primitive::os_error(errno, process); } +// POSIX writes have already completed when the write primitive returns. +PRIMITIVE(write_result) { + ARGS(IntResource, fd_resource, int, written); + USE(fd_resource); + return Smi::from(written); +} + PRIMITIVE(fd) { ARGS(IntResource, fd_resource); int fd = fd_resource->id(); diff --git a/src/resources/pipe_win.cc b/src/resources/pipe_win.cc index 9f2169c76..63a675c4e 100644 --- a/src/resources/pipe_win.cc +++ b/src/resources/pipe_win.cc @@ -72,24 +72,25 @@ class HandlePipeResource : public WindowsResource { TAG(PipeResource); HandlePipeResource(ResourceGroup* resource_group, HANDLE handle, HANDLE event) : WindowsResource(resource_group), handle_(handle) { - overlapped_.hEvent = event; + overlapped_.set_event(event); } HANDLE handle() { return handle_; } std::vector events() override { - return std::vector( { overlapped_.hEvent } ); + return std::vector( { overlapped_.event() } ); } void do_close() override { - CloseHandle(overlapped_.hEvent); + overlapped_.cancel_and_wait(handle_); CloseHandle(handle_); + CloseHandle(overlapped_.event()); } - OVERLAPPED* overlapped() { return &overlapped_; } + WindowsOverlapped& overlapped() { return overlapped_; } private: HANDLE handle_; - OVERLAPPED overlapped_{}; + WindowsOverlapped overlapped_; }; class ReadPipeResource : public HandlePipeResource { @@ -107,15 +108,12 @@ class ReadPipeResource : public HandlePipeResource { bool issue_read_request() { read_ready_ = false; read_count_ = 0; - bool success = ReadFile(handle(), read_data_, READ_BUFFER_SIZE, &read_count_, overlapped()); - if (!success && WSAGetLastError() != ERROR_IO_PENDING) { - return false; - } - return true; + bool success = ReadFile(handle(), read_data_, READ_BUFFER_SIZE, &read_count_, overlapped().get()); + return overlapped().issued(success, GetLastError()); } bool receive_read_response() { - bool overlapped_result = GetOverlappedResult(handle(), overlapped(), &read_count_, false); + bool overlapped_result = GetOverlappedResult(handle(), overlapped().get(), &read_count_, false); return overlapped_result; } @@ -150,7 +148,8 @@ class WritePipeResource : public HandlePipeResource { bool ready_for_write() const { return write_ready_; } - + // The write primitive preserves its queued-count result for older host + // packages. New callers must use write_result before reporting success. bool send(const uint8* buffer, word length) { if (write_buffer_ != null) free(write_buffer_); @@ -158,14 +157,16 @@ class WritePipeResource : public HandlePipeResource { // We need to copy the buffer out to a long-lived heap object. write_buffer_ = static_cast(malloc(length)); + if (write_buffer_ == null) { + write_ready_ = true; + SetLastError(ERROR_NOT_ENOUGH_MEMORY); + return false; + } memcpy(write_buffer_, buffer, length); DWORD tmp; - bool send_result = WriteFile(handle(), write_buffer_, length, &tmp, overlapped()); - if (!send_result && WSAGetLastError() != ERROR_IO_PENDING) { - return false; - } - return true; + bool send_result = WriteFile(handle(), write_buffer_, length, &tmp, overlapped().get()); + return overlapped().issued(send_result, GetLastError()); } private: @@ -462,6 +463,19 @@ PRIMITIVE(write) { return Smi::from(to - from); } +// Returns null while a queued write is pending. Waiting happens in Toit, so +// neither the scheduler thread nor the event thread blocks on a slow reader. +PRIMITIVE(write_result) { + ARGS(WritePipeResource, pipe_resource, int, written); + USE(written); + DWORD count; + if (!GetOverlappedResult(pipe_resource->handle(), pipe_resource->overlapped().get(), &count, FALSE)) { + if (GetLastError() == ERROR_IO_INCOMPLETE) return process->null_object(); + WINDOWS_ERROR; + } + return Smi::from(count); +} + PRIMITIVE(read) { ARGS(ReadPipeResource, read_resource); diff --git a/src/resources/tcp_win.cc b/src/resources/tcp_win.cc index f1b772500..f22639ef7 100644 --- a/src/resources/tcp_win.cc +++ b/src/resources/tcp_win.cc @@ -80,13 +80,12 @@ class TcpSocketResource : public TcpSocketBaseResource { public: TAG(TcpSocketResource); TcpSocketResource(TcpResourceGroup* resource_group, SOCKET socket, - HANDLE read_event, HANDLE write_event, HANDLE auxiliary_event) + HANDLE read_event, HANDLE auxiliary_event) : TcpSocketBaseResource(resource_group, socket) , auxiliary_event_(auxiliary_event) { read_buffer_.buf = read_data_; read_buffer_.len = READ_BUFFER_SIZE; - read_overlapped_.hEvent = read_event; - write_overlapped_.hEvent = write_event; + read_overlapped_.set_event(read_event); if (!issue_read_request()) { int error_code = WSAGetLastError(); if (error_code == WSAECONNRESET) { @@ -100,37 +99,36 @@ class TcpSocketResource : public TcpSocketBaseResource { } } - ~TcpSocketResource() override { - if (write_buffer_.buf != null) free(write_buffer_.buf); - } - DWORD read_count() const { return read_count_; } char* read_buffer() const { return read_buffer_.buf; } - bool ready_for_write() const { return write_ready_; } bool ready_for_read() const { return read_ready_; } bool closed() const { return closed_; } std::vector events() override { return std::vector({ - read_overlapped_.hEvent, - write_overlapped_.hEvent, + read_overlapped_.event(), auxiliary_event_ }); } uint32_t on_event(HANDLE event, uint32_t state) override { - if (event == read_overlapped_.hEvent) { + if (event == read_overlapped_.event()) { read_ready_ = true; state |= TCP_READ; - } else if (event == write_overlapped_.hEvent) { - write_ready_ = true; - state |= TCP_WRITE; } else if (event == auxiliary_event_) { WSANETWORKEVENTS network_events; if (WSAEnumNetworkEvents(socket(), NULL, &network_events) == SOCKET_ERROR) { set_error_code(WSAGetLastError()); - state |= TCP_ERROR; - }; + return state | TCP_ERROR; + } + if (network_events.lNetworkEvents & FD_WRITE) { + if (network_events.iErrorCode[FD_WRITE_BIT] == 0) { + state |= TCP_WRITE; + } else { + set_error_code(network_events.iErrorCode[FD_WRITE_BIT]); + state |= TCP_ERROR; + } + } if (network_events.lNetworkEvents & FD_CLOSE) { if (network_events.iErrorCode[FD_CLOSE_BIT] == 0) { state |= TCP_READ; @@ -150,65 +148,47 @@ class TcpSocketResource : public TcpSocketBaseResource { } void do_close() override { + // Reaping the outstanding read also matters for the peer: closing a socket + // that has a pending overlapped WSARecv is an abortive close, so the peer + // gets an RST and loses the data it has not read yet. + read_overlapped_.cancel_and_wait(reinterpret_cast(socket())); TcpSocketBaseResource::do_close(); - CloseHandle(read_overlapped_.hEvent); - CloseHandle(write_overlapped_.hEvent); + CloseHandle(read_overlapped_.event()); + CloseHandle(auxiliary_event_); } bool issue_read_request() { read_ready_ = false; read_count_ = 0; DWORD flags = 0; - int receive_result = WSARecv(socket(), &read_buffer_, 1, NULL, &flags, &read_overlapped_, NULL); - if (receive_result == SOCKET_ERROR && WSAGetLastError() != WSA_IO_PENDING) { - return false; - } - return true; + int receive_result = WSARecv(socket(), &read_buffer_, 1, NULL, &flags, read_overlapped_.get(), NULL); + return read_overlapped_.issued(receive_result == 0, WSAGetLastError()); } bool receive_read_response() { DWORD flags; - bool overlapped_result = WSAGetOverlappedResult(socket(), &read_overlapped_, &read_count_, false, &flags); + bool overlapped_result = WSAGetOverlappedResult(socket(), read_overlapped_.get(), &read_count_, false, &flags); if (read_count_ == 0) closed_ = true; return overlapped_result; } - bool send(const uint8* buffer, word length) { - if (write_buffer_.buf != null) { - free(write_buffer_.buf); - write_buffer_.buf = null; - } - - write_ready_ = false; - - // We need to copy the buffer out to a long-lived heap object - write_buffer_.buf = static_cast(malloc(length)); - if (!write_buffer_.buf) { - WSASetLastError(ERROR_NOT_ENOUGH_MEMORY); - return false; - } - memcpy(write_buffer_.buf, buffer, length); - write_buffer_.len = length; - - int send_result = WSASend(socket(), &write_buffer_, 1, NULL, 0, &write_overlapped_, NULL); - - if (send_result == SOCKET_ERROR && WSAGetLastError() != WSA_IO_PENDING) { - return false; - } - return true; + int send(const uint8* buffer, int length) { + // WSAEventSelect puts the socket in nonblocking mode. Only report bytes + // accepted by the transport, so close cannot cancel a queued WSASend. + // FD_WRITE wakes the writer once a full send buffer has room again. + // Always retry the nonblocking send. Caching readiness here can lose an + // FD_WRITE handled between send returning WSAEWOULDBLOCK and clearing the + // cached flag, leaving the writer asleep on a writable socket. + return ::send(socket(), reinterpret_cast(buffer), length, 0); } private: WSABUF read_buffer_{}; char read_data_[READ_BUFFER_SIZE]{}; - OVERLAPPED read_overlapped_{}; + WindowsOverlapped read_overlapped_; DWORD read_count_ = 0; bool read_ready_ = false; - WSABUF write_buffer_{}; - OVERLAPPED write_overlapped_{}; - bool write_ready_ = true; - HANDLE auxiliary_event_; bool closed_ = false; }; @@ -256,13 +236,13 @@ PRIMITIVE(init) { } static Object* create_events(Process* process, SOCKET socket, WSAEVENT& read_event, - WSAEVENT& write_event, WSAEVENT& auxiliary_event) { + WSAEVENT& auxiliary_event) { auxiliary_event = WSACreateEvent(); if (auxiliary_event == WSA_INVALID_EVENT) { WINDOWS_ERROR; } - if (WSAEventSelect(socket, auxiliary_event, FD_CLOSE) == SOCKET_ERROR) { + if (WSAEventSelect(socket, auxiliary_event, FD_WRITE | FD_CLOSE) == SOCKET_ERROR) { close_handle_keep_errno(auxiliary_event); WINDOWS_ERROR; } @@ -273,13 +253,6 @@ static Object* create_events(Process* process, SOCKET socket, WSAEVENT& read_eve WINDOWS_ERROR; } - write_event = WSACreateEvent(); - if (write_event == WSA_INVALID_EVENT) { - close_handle_keep_errno(read_event); - close_handle_keep_errno(auxiliary_event); - WINDOWS_ERROR; - } - return null; } @@ -299,25 +272,24 @@ PRIMITIVE(connect) { } ToitSocketAddress socket_address(address.address(), address.length(), port); - int result = connect(socket, socket_address.as_socket_address(), socket_address.port()); + int result = connect(socket, socket_address.as_socket_address(), socket_address.size()); if (result == SOCKET_ERROR && WSAGetLastError() != WSAEINPROGRESS) { close_keep_errno(socket); WINDOWS_ERROR; } - WSAEVENT read_event, write_event, auxiliary_event; - auto error = create_events(process, socket, read_event, write_event, auxiliary_event); + WSAEVENT read_event, auxiliary_event; + auto error = create_events(process, socket, read_event, auxiliary_event); if (error) { close_keep_errno(socket); return error; } - auto tcp_resource = _new TcpSocketResource(resource_group, socket, read_event, write_event, auxiliary_event); + auto tcp_resource = _new TcpSocketResource(resource_group, socket, read_event, auxiliary_event); if (!tcp_resource) { close_keep_errno(socket); close_handle_keep_errno(read_event); - close_handle_keep_errno(write_event); close_handle_keep_errno(auxiliary_event); FAIL(MALLOC_FAILED); } @@ -342,19 +314,19 @@ PRIMITIVE(accept) { WINDOWS_ERROR; } - WSAEVENT read_event, write_event, auxiliary_event; - auto error = create_events(process, socket, read_event, write_event, auxiliary_event); + WSAEVENT read_event, auxiliary_event; + auto error = create_events(process, socket, read_event, auxiliary_event); if (error) { close_keep_errno(socket); return error; } - auto tcp_resource = _new TcpSocketResource(resource_group, socket, read_event, write_event, auxiliary_event); + auto tcp_resource = _new TcpSocketResource(resource_group, socket, read_event, auxiliary_event); if (!tcp_resource) { close_keep_errno(socket); close_handle_keep_errno(read_event); - close_handle_keep_errno(write_event); + close_handle_keep_errno(auxiliary_event); FAIL(MALLOC_FAILED); } @@ -424,11 +396,12 @@ PRIMITIVE(write) { if (from < 0 || from > to || to > data.length()) FAIL(OUT_OF_BOUNDS); - if (!tcp_resource->ready_for_write()) return Smi::from(-1); - - if (!tcp_resource->send(data.address() + from, to - from)) WINDOWS_ERROR; - - return Smi::from(to-from); + int written = tcp_resource->send(data.address() + from, to - from); + if (written == SOCKET_ERROR) { + if (WSAGetLastError() == WSAEWOULDBLOCK) return Smi::from(-1); + WINDOWS_ERROR; + } + return Smi::from(written); } PRIMITIVE(read) { diff --git a/src/resources/uart_win.cc b/src/resources/uart_win.cc index 604d42434..de03b968a 100644 --- a/src/resources/uart_win.cc +++ b/src/resources/uart_win.cc @@ -43,9 +43,9 @@ class UartResource : public WindowsResource { UartResource(ResourceGroup* group, HANDLE uart, HANDLE read_event, HANDLE write_event, HANDLE error_event) : WindowsResource(group) , uart_(uart) { - read_overlapped_.hEvent = read_event; - write_overlapped_.hEvent = write_event; - comm_events_overlapped_.hEvent = error_event; + read_overlapped_.set_event(read_event); + write_overlapped_.set_event(write_event); + comm_events_overlapped_.set_event(error_event); set_state(kWriteState); @@ -75,42 +75,42 @@ class UartResource : public WindowsResource { DWORD error_code() const { return error_code_; } void do_close() override { - CloseHandle(read_overlapped_.hEvent); - CloseHandle(write_overlapped_.hEvent); - CloseHandle(comm_events_overlapped_.hEvent); + read_overlapped_.cancel_and_wait(uart_); + write_overlapped_.cancel_and_wait(uart_); + // Also keeps the completion from writing to event_mask_. + comm_events_overlapped_.cancel_and_wait(uart_); CloseHandle(uart_); + CloseHandle(read_overlapped_.event()); + CloseHandle(write_overlapped_.event()); + CloseHandle(comm_events_overlapped_.event()); } std::vector events() override { return std::vector({ - read_overlapped_.hEvent, - write_overlapped_.hEvent, - comm_events_overlapped_.hEvent + read_overlapped_.event(), + write_overlapped_.event(), + comm_events_overlapped_.event() }); } bool issue_comm_events_request() { - bool succeeded = WaitCommEvent(uart_, &event_mask_, &comm_events_overlapped_); - if (!succeeded && GetLastError() != ERROR_IO_PENDING) { - return false; - } - return true; + bool succeeded = WaitCommEvent(uart_, &event_mask_, comm_events_overlapped_.get()); + return comm_events_overlapped_.issued(succeeded, GetLastError()); } bool issue_read_request() { read_ready_ = false; read_count_ = 0; - bool success = ReadFile(uart_, read_data_, READ_BUFFER_SIZE, &read_count_, &read_overlapped_); - if (!success && GetLastError() != ERROR_IO_PENDING) { - return false; - } - return true; + bool success = ReadFile(uart_, read_data_, READ_BUFFER_SIZE, &read_count_, read_overlapped_.get()); + return read_overlapped_.issued(success, GetLastError()); } bool receive_read_response() { - return GetOverlappedResult(uart_, &read_overlapped_, &read_count_, false); + return GetOverlappedResult(uart_, read_overlapped_.get(), &read_count_, false); } + // Reports success as soon as the write is queued. Closing the port cancels + // a write that is still pending, so its data can be lost. bool send(const uint8* buffer, word length) { // The caller must ensure that the previous write has completed before // calling send again. If write_ready_ is false, the previous WriteFile @@ -125,24 +125,20 @@ class UartResource : public WindowsResource { memcpy(write_buffer_, buffer, length); DWORD tmp; - bool send_result = WriteFile(uart_, write_buffer_, length, &tmp, &write_overlapped_); - if (!send_result && GetLastError() != ERROR_IO_PENDING) { - return false; - } - - return true; + bool send_result = WriteFile(uart_, write_buffer_, length, &tmp, write_overlapped_.get()); + return write_overlapped_.issued(send_result, GetLastError()); } uint32_t on_event(HANDLE event, uint32_t state) override { - if (event == read_overlapped_.hEvent) { + if (event == read_overlapped_.event()) { read_ready_ = true; state |= kReadState; - } else if (event == write_overlapped_.hEvent) { + } else if (event == write_overlapped_.event()) { write_ready_ = true; state |= kWriteState; - } else if (event == comm_events_overlapped_.hEvent) { + } else if (event == comm_events_overlapped_.event()) { DWORD tmp; - bool succeeded = GetOverlappedResult(uart_, &comm_events_overlapped_, &tmp, false); + bool succeeded = GetOverlappedResult(uart_, comm_events_overlapped_.get(), &tmp, false); if (!succeeded) { error_code_ = GetLastError(); } else { @@ -170,15 +166,15 @@ class UartResource : public WindowsResource { bool dtr_ = false; char read_data_[READ_BUFFER_SIZE]{}; - OVERLAPPED read_overlapped_{}; + WindowsOverlapped read_overlapped_; DWORD read_count_ = 0; bool read_ready_ = false; - OVERLAPPED write_overlapped_{}; + WindowsOverlapped write_overlapped_; char* write_buffer_ = null; bool write_ready_ = true; - OVERLAPPED comm_events_overlapped_{}; + WindowsOverlapped comm_events_overlapped_; DWORD event_mask_ = 0; DWORD error_code_ = ERROR_SUCCESS; diff --git a/src/resources/udp_win.cc b/src/resources/udp_win.cc index e5e5e3058..ce525057f 100644 --- a/src/resources/udp_win.cc +++ b/src/resources/udp_win.cc @@ -63,8 +63,8 @@ class UdpSocketResource : public WindowsResource { , socket_(socket) { read_buffer_.buf = read_data_; read_buffer_.len = READ_BUFFER_SIZE; - read_overlapped_.hEvent = read_event; - write_overlapped_.hEvent = write_event; + read_overlapped_.set_event(read_event); + write_overlapped_.set_event(write_event); set_state(UDP_WRITE); } @@ -91,14 +91,14 @@ class UdpSocketResource : public WindowsResource { bool ready_for_write() const { return write_ready_; } std::vector events() override { - return std::vector({read_overlapped_.hEvent, write_overlapped_.hEvent }); + return std::vector({read_overlapped_.event(), write_overlapped_.event() }); } uint32_t on_event(HANDLE event, uint32_t state) override { - if (event == read_overlapped_.hEvent) { + if (event == read_overlapped_.event()) { read_ready_ = true; state |= UDP_READ; - } else if (event == write_overlapped_.hEvent) { + } else if (event == write_overlapped_.event()) { write_ready_ = true; state |= UDP_WRITE; } @@ -107,9 +107,11 @@ class UdpSocketResource : public WindowsResource { } void do_close() override { + read_overlapped_.cancel_and_wait(reinterpret_cast(socket_)); + write_overlapped_.cancel_and_wait(reinterpret_cast(socket_)); closesocket(socket_); - CloseHandle(read_overlapped_.hEvent); - CloseHandle(write_overlapped_.hEvent); + CloseHandle(read_overlapped_.event()); + CloseHandle(write_overlapped_.event()); } bool issue_read_request() { @@ -119,16 +121,13 @@ class UdpSocketResource : public WindowsResource { int receive_result = WSARecvFrom(socket_, &read_buffer_, 1, NULL, &flags, read_peer_address_.as_socket_address(), read_peer_address_.size_pointer(), - &read_overlapped_, NULL); - if (receive_result == SOCKET_ERROR && WSAGetLastError() != WSA_IO_PENDING) { - return false; - } - return true; + read_overlapped_.get(), NULL); + return read_overlapped_.issued(receive_result == 0, WSAGetLastError()); } bool receive_read_response() { DWORD flags; - bool overlapped_result = WSAGetOverlappedResult(socket_, &read_overlapped_, &read_count_, false, &flags); + bool overlapped_result = WSAGetOverlappedResult(socket_, read_overlapped_.get(), &read_count_, false, &flags); return overlapped_result; } @@ -151,12 +150,12 @@ class UdpSocketResource : public WindowsResource { send_result = WSASendTo(socket_, &write_buffer_, 1, &tmp, 0, socket_address->as_socket_address(), socket_address->size(), - &write_overlapped_, NULL); + write_overlapped_.get(), NULL); } else { - send_result = WSASend(socket_, &write_buffer_, 1, &tmp, 0, &write_overlapped_, NULL); + send_result = WSASend(socket_, &write_buffer_, 1, &tmp, 0, write_overlapped_.get(), NULL); } - if (send_result == SOCKET_ERROR && WSAGetLastError() != WSA_IO_PENDING) { + if (!write_overlapped_.issued(send_result == 0, WSAGetLastError())) { return false; } @@ -168,13 +167,13 @@ class UdpSocketResource : public WindowsResource { WSABUF read_buffer_{}; char read_data_[READ_BUFFER_SIZE]{}; - OVERLAPPED read_overlapped_{}; + WindowsOverlapped read_overlapped_; DWORD read_count_ = 0; ToitSocketAddress read_peer_address_; bool read_ready_ = false; WSABUF write_buffer_{}; - OVERLAPPED write_overlapped_{}; + WindowsOverlapped write_overlapped_; bool write_ready_ = true; DWORD error_code_ = ERROR_SUCCESS; diff --git a/tests/pipe-close-pending-write-test-compiler.toit b/tests/pipe-close-pending-write-test-compiler.toit new file mode 100644 index 000000000..5cbc6a0d7 --- /dev/null +++ b/tests/pipe-close-pending-write-test-compiler.toit @@ -0,0 +1,62 @@ +// Copyright (C) 2026 Toit contributors. +// Use of this source code is governed by a Zero-Clause BSD license that can +// be found in the tests/LICENSE file. + +import expect show * +import host.pipe +import monitor show * +import system + +// Closing a pipe while a write is still pending must not let the write's +// completion land in the freed pipe resource. On Windows the completion +// arrived when the reader went away, and wrote into whatever had reused the +// resource's memory. Here that is one of many freshly allocated external byte +// arrays of the same size. + +// Bigger than the 8192-byte pipe buffer on Windows, so the write stays +// pending. Smaller than the pipe buffers elsewhere, so the write doesn't block. +PAYLOAD-SIZE ::= 12_000 +// Covers the size of the Windows pipe resource. +SPRAY-SIZES ::= [104, 112, 120, 128, 136] +SPRAY-COUNT ::= 500 +ITERATIONS ::= 5 + +// Marked as compiler test so we get the toit.run path. +main args: + if args.size == 1 and args[0] == "CHILD": + // Never read stdin; wait to be killed. + sleep --ms=60_000 + return + + toit-run := args[0] + ITERATIONS.repeat: test toit-run + +test toit-run/string: + process := pipe.fork + --create-stdin + toit-run + [toit-run, system.program-path, "CHILD"] + + started := Latch + finished := Latch + task:: + started.set true + error := catch: process.stdin.out.write (ByteArray PAYLOAD-SIZE) + finished.set (error or false) + started.get + // Let the writer fill the pipe and suspend before closing it from this task. + sleep --ms=50 + process.stdin.close + finished.get + + spray := [] + SPRAY-SIZES.do: | size | + SPRAY-COUNT.repeat: spray.add (ByteArray.external size) + + // The reader goes away, which completes a write that is still pending. + process.kill --hard + process.wait + sleep --ms=50 + + spray.do: | bytes/ByteArray | + bytes.do: expect-equals 0 it diff --git a/tests/tcp-write-close-test.toit b/tests/tcp-write-close-test.toit new file mode 100644 index 000000000..28a2b5fb0 --- /dev/null +++ b/tests/tcp-write-close-test.toit @@ -0,0 +1,83 @@ +// Copyright (C) 2026 Toit contributors. +// Use of this source code is governed by a Zero-Clause BSD license that can +// be found in the tests/LICENSE file. + +import expect show * +import monitor show * +import net +import net.modules.tcp + +// Data written just before a close must reach the peer, even when the peer +// only starts reading after the close has happened. On Windows a close used +// to be abortive, which sent an RST and discarded the payload. + +HEADER ::= "4096\n".to-byte-array +PAYLOAD ::= ByteArray 4096: it & 0xff + +main: + test-small-write + test-backpressure + +test-small-write: + network := net.open + port := Latch + closed := Latch + + task:: + server := tcp.TcpServerSocket network + server.listen "127.0.0.1" 0 + port.set server.local-address.port + socket := server.accept + // Two separate writes, then an immediate close. + socket.out.write HEADER + socket.out.write PAYLOAD + socket.close + server.close + closed.set true + + socket := tcp.TcpSocket network + socket.connect "127.0.0.1" port.get + // Don't read anything until the peer has written and closed. The data is + // then sitting in our receive buffer, which is what an abortive close on + // the other side throws away. + closed.get + sleep --ms=100 + received := socket.in.read-all + socket.close + expect-equals HEADER.size + PAYLOAD.size received.size + expect-equals HEADER + PAYLOAD received + +// Force the writer to run out of transport buffer space. A second task must +// remain runnable while the peer is not reading, and every byte must arrive. +// Repeat on the same socket so readiness must recover after each stall. +test-backpressure: + network := net.open + port := Latch + started := Channel 1 + finished := Latch + rounds := 4 + payload := ByteArray (16 * 1024 * 1024): it & 0xff + + task:: + server := tcp.TcpServerSocket network + server.listen "127.0.0.1" 0 + port.set server.local-address.port + socket := server.accept + rounds.repeat: + started.send true + socket.out.write payload + socket.close + server.close + finished.set true + + socket := tcp.TcpSocket network + socket.connect "127.0.0.1" port.get + reader := socket.in + rounds.repeat: + started.receive + sleep --ms=100 + received := reader.read-bytes payload.size + expect-equals payload received + expect-null reader.read + socket.close + finished.get diff --git a/tests/zlib-gzip-test.toit b/tests/zlib-gzip-test.toit index f1a5bb393..380513932 100644 --- a/tests/zlib-gzip-test.toit +++ b/tests/zlib-gzip-test.toit @@ -4,6 +4,8 @@ import expect show * import zlib +import host.directory +import host.file import host.pipe import io @@ -34,22 +36,24 @@ test-gzip str/string: expect-equals str result.bytes.to-string gzip-compress data/ByteArray -> ByteArray: - // Use host 'gzip' to compress. - proc := pipe.fork - --use-path - --create-stdin - --create-stdout - "gzip" - ["gzip", "-c"] - - writer := proc.stdin.out - writer.write data - writer.close - - result := io.Buffer - reader := proc.stdout.in - while chunk := reader.read: - result.write chunk - - proc.wait - return result.bytes + // Use host 'gzip' to compress. Give it a file rather than writing to its + // stdin: on Windows, closing a pipe cancels a write that is still pending. + tmp-dir := directory.mkdtemp "/tmp/zlib-gzip-test-" + try: + path := "$tmp-dir/data" + file.write-contents --path=path data + proc := pipe.fork + --use-path + --create-stdout + "gzip" + ["gzip", "-c", path] + + result := io.Buffer + reader := proc.stdout.in + while chunk := reader.read: + result.write chunk + + proc.wait + return result.bytes + finally: + directory.rmdir --recursive tmp-dir