2#include "pedigree/kernel/process/Readiness.h"
3#include "pedigree/kernel/process/Semaphore.h"
4#include "pedigree/kernel/process/TerminationDeferral.h"
5#include "pedigree/kernel/process/Thread.h"
6#include "pedigree/kernel/processor/Processor.h"
7#include "pedigree/kernel/processor/ProcessorInformation.h"
8#include "pedigree/kernel/syscallError.h"
9#include "pedigree/kernel/time/Time.h"
14#include "FileDescriptor.h"
15#include "PosixSubsystem.h"
16#include "net-syscalls.h"
17#include "poll-syscalls.h"
18#include "recvmmsg-syscalls.h"
25 m_Semaphore->release();
38 return deadline.type == PollDeadlineType::Immediate ||
39 (deadline.type == PollDeadlineType::Finite && Time::getTicks() >= deadline.expires);
54 ReadyRead | ReadyError | ReadyHangup | ReadyReadHangup | ReadyInvalid,
55 m_Observer, m_Subscription);
61 (ReadyRead | ReadyError | ReadyHangup | ReadyReadHangup | ReadyInvalid))
63 if (thread.getInterruptionReason() == Thread::InterruptedBySignal)
65 if (deadlineExpired(deadline))
68 size_t seconds = 0, microseconds = 0;
69 if (deadline.type == PollDeadlineType::Finite) {
70 const uint64_t now = Time::getTicks();
71 if (now >= deadline.expires)
73 const uint64_t remaining = deadline.expires - now;
74 const uint64_t rounded = remaining / Time::Multiplier::Microsecond +
75 (remaining % Time::Multiplier::Microsecond != 0);
76 seconds =
static_cast<size_t>(rounded / 1000000);
77 microseconds =
static_cast<size_t>(rounded % 1000000);
79 Semaphore::SemaphoreError error = Semaphore::NoError;
80 const bool signalled = m_Semaphore->acquireWithError(1, seconds, microseconds, error);
82 while (m_Semaphore->tryAcquire()) {
84 }
else if (error != Semaphore::TimedOut) {
87 (ReadyRead | ReadyError | ReadyHangup | ReadyReadHangup | ReadyInvalid)
100ReceiveResult receiveMessages(
int fd,
LinuxMmsghdr* messages,
unsigned int count,
104 thread.clearInterruption();
109 if (duration.tv_sec < 0 || duration.tv_nsec < 0 ||
110 duration.tv_nsec >=
static_cast<int64_t
>(Time::Multiplier::Second))
113 const PollDeadline deadline = posix_poll_deadline(timeout ? &duration : nullptr);
116 if (!subsystem || !subsystem->acquireFileDescriptor(fd, descriptor))
119 return {-1, ENOTSOCK};
121 const int deferred = socket.takeReceiveError();
123 return {-1, deferred};
126 if (count >
static_cast<unsigned int>(INT_MAX))
130 unsigned int received = 0;
132 bool interrupted =
false;
133 uintptr_t address =
reinterpret_cast<uintptr_t
>(messages);
134 while (received < count) {
140 const int receiveFlags =
static_cast<int>(flags & ~MSG_WAITFORONE) | MSG_DONTWAIT;
141 interrupted |= thread.getInterruptionReason() == Thread::InterruptedBySignal;
143 const ssize_t bytes = posix_recvmsg_user_descriptor(descriptor, &
message->msg_hdr, receiveFlags,
148 if (timeout && deadlineExpired(deadline))
161 if ((flags & MSG_DONTWAIT) || !socket.isBlocking() || (received && (flags & MSG_WAITFORONE)))
163 if (deadlineExpired(deadline)) {
167 const ReadyMask readyState = socket.
queryReady(
true,
false);
168 if (readyState & ReadyInvalid) {
172 if (readyState & ReadyError) {
174 socklen_t length =
sizeof(socketError);
175 if (socket.getsockopt(SOL_SOCKET, SO_ERROR, &socketError, &length) < 0) {
184 if (!wait.prepare(socket)) {
188 const int ready = wait.wait(socket, deadline, thread);
190 error = ready ? EINTR : 0;
195 if (received && error && error != EAGAIN)
196 socket.deferReceiveError(error);
197 if (timeout && received) {
199 if (deadline.type == PollDeadlineType::Finite) {
200 const uint64_t now = Time::getTicks();
201 if (now < deadline.expires) {
202 const uint64_t nanoseconds = deadline.expires - now;
203 remaining.tv_sec = nanoseconds / Time::Multiplier::Second;
204 remaining.tv_nsec = nanoseconds % Time::Multiplier::Second;
210 if (received || !error) {
211 thread.clearInterruption();
212 return {
static_cast<int>(received), 0};
218int posix_recvmmsg(
int fd,
LinuxMmsghdr* messages,
unsigned int count,
unsigned int flags,
220 if (flags & 0x80000000U) {
221 syscallError(EINVAL);
224 const ReceiveResult result = receiveMessages(fd, messages, count, flags, timeout);
225 syscallError(result.error);
SharedPointer< NetworkSyscalls > networkImpl
Network syscall implementation for this descriptor (if it's a socket).
virtual ReadyMask queryReady(bool reading, bool writing)
static bool copyFromUser(void *destination, const void *source, size_t count, size_t elementSize=1)
static bool copyToUser(void *destination, const void *source, size_t count, size_t elementSize=1)
static ProcessorInformation & information()
virtual void readinessChanged(ReadyMask mask)=0
MUST_USE_RESULT bool subscribeReadiness(ReadyMask interest, const SharedPointer< ReadinessObserver > &observer, ReadinessSubscription &subscription)
static SharedPointer< T > tryAdopt(T *ptr)
static SharedPointer< T > tryAllocate(Args...)
Process * getParent() const