Skip to content
Merged
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
1 change: 1 addition & 0 deletions src/compiler/propagation/type_primitive_pipe.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
10 changes: 10 additions & 0 deletions src/event_sources/event_win.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
31 changes: 31 additions & 0 deletions src/event_sources/event_win.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
1 change: 1 addition & 0 deletions src/primitive.h
Original file line number Diff line number Diff line change
Expand Up @@ -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) \
Expand Down
7 changes: 7 additions & 0 deletions src/resources/pipe_posix.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
48 changes: 31 additions & 17 deletions src/resources/pipe_win.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<HANDLE> events() override {
return std::vector<HANDLE>( { overlapped_.hEvent } );
return std::vector<HANDLE>( { overlapped_.event() } );
}

void do_close() override {
CloseHandle(overlapped_.hEvent);
overlapped_.cancel_and_wait(handle_);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
CloseHandle(handle_);
CloseHandle(overlapped_.event());
}

OVERLAPPED* overlapped() { return &overlapped_; }
WindowsOverlapped& overlapped() { return overlapped_; }
private:
HANDLE handle_;
OVERLAPPED overlapped_{};
WindowsOverlapped overlapped_;
};

class ReadPipeResource : public HandlePipeResource {
Expand All @@ -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;
}

Expand Down Expand Up @@ -150,22 +148,25 @@ 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_);

write_ready_ = false;

// We need to copy the buffer out to a long-lived heap object.
write_buffer_ = static_cast<char*>(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:
Expand Down Expand Up @@ -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);

Expand Down
Loading
Loading