20#include "UnixFilesystem.h"
21#include "pedigree/kernel/LockGuard.h"
22#include "pedigree/kernel/process/Mutex.h"
23#include "pedigree/kernel/process/Process.h"
24#include "pedigree/kernel/process/Thread.h"
25#include "pedigree/kernel/processor/Processor.h"
26#include "pedigree/kernel/syscallError.h"
30#include "modules/subsys/posix/FileDescriptor.h"
31#include "modules/subsys/posix/logging.h"
32#include "modules/system/vfs/VFS.h"
34String UnixFilesystem::m_VolumeLabel(
"unix");
35Mutex UnixFilesystem::m_NamespaceLock;
36Mutex UnixSocket::m_ConnectionLock;
37Mutex SocketRights::m_InFlightLock;
38size_t SocketRights::m_InFlight = 0;
40SocketRights::SocketRights(
size_t reservation)
41 : m_Descriptors(reservation), m_Reservation(reservation) {}
43SocketRights::~SocketRights() {
44 for (
auto descriptor : m_Descriptors) {
47 m_Descriptors.
clear(
true);
50 assert(m_InFlight >= m_Reservation);
51 m_InFlight -= m_Reservation;
56 if (!descriptorCount || descriptorCount > MaximumDescriptors) {
62 if (descriptorCount > MaximumInFlight - m_InFlight) {
65 m_InFlight += descriptorCount;
74 assert(m_Descriptors.
count() < m_Reservation);
78size_t SocketRights::count()
const {
79 return m_Descriptors.
count();
83 return m_Descriptors[index];
86#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
87size_t SocketRights::inFlightForTest() {
94#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
95UnixStreamControlLockHook g_UnixStreamControlLockHook =
nullptr;
97void invokeUnixStreamControlLockHook() {
98 UnixStreamControlLockHook hook =
99 __atomic_exchange_n(&g_UnixStreamControlLockHook,
100 static_cast<UnixStreamControlLockHook
>(
nullptr), __ATOMIC_ACQ_REL);
107bool currentThreadWasInterrupted() {
108#if defined(PEDIGREE_EXTERNAL_SOURCE)
112 return thread && thread->getInterruptionReason() == Thread::InterruptedBySignal;
116void captureStreamInterruption(
bool block,
bool* interrupted) {
117 if (block && interrupted && currentThreadWasInterrupted()) {
122enum class StreamSerializationWait {
127class StreamSerializationGuard {
129#if defined(PEDIGREE_EXTERNAL_SOURCE)
131 : m_Mutex(mutex), m_Acquired(m_Mutex.acquire()) {}
133 StreamSerializationGuard(
Semaphore& semaphore, StreamSerializationWait wait)
134 : m_TerminationDeferral(true), m_Semaphore(semaphore), m_Acquired(false) {
135 if (wait == StreamSerializationWait::Nonblocking) {
136 m_Acquired = m_Semaphore.tryAcquire();
138 Semaphore::SemaphoreError error = Semaphore::NoError;
139 m_Acquired = m_Semaphore.acquireWithError(1, 0, 0, error);
148 ~StreamSerializationGuard() {
150#if defined(PEDIGREE_EXTERNAL_SOURCE)
153 m_Semaphore.release();
158 explicit operator bool()
const {
163#if defined(PEDIGREE_EXTERNAL_SOURCE)
173#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
174void setUnixStreamControlLockHookForTest(UnixStreamControlLockHook hook) {
175 __atomic_store_n(&g_UnixStreamControlLockHook, hook, __ATOMIC_RELEASE);
179UnixSocketConnection::Stream::ControlGuard::ControlGuard(Stream& stream)
180 : m_Stream(stream), m_Guard(stream.m_ControlLock) {
181#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
182 invokeUnixStreamControlLockHook();
186UnixSocketConnection::Stream::ControlGuard::~ControlGuard() {
187 m_Stream.discardControlsIfRequested();
188 m_Stream.m_ControlLock.
release();
194 if (m_Stream.m_DiscardControlsRequested) {
195 ControlGuard cleanupGuard(m_Stream);
199UnixSocketConnection::Stream::Stream(
bool packets)
200 : m_Packets(packets),
201 m_Records(MAX_UNIX_PACKET_BACKLOG * sizeof(uintptr_t)),
203 m_Bytes(MAX_UNIX_STREAM_QUEUE),
204#if defined(PEDIGREE_EXTERNAL_SOURCE)
209 m_ReceiveLock(1, true),
212 m_DiscardControlsRequested(false),
218UnixSocketConnection::Stream::~Stream() {
221 m_DiscardControlsRequested =
false;
224size_t UnixSocketConnection::Stream::write(
const uint8_t* buffer,
size_t count,
bool block,
227 struct iovec vector = {
const_cast<uint8_t*
>(buffer), count};
228 return writeVectors(&vector, 1, block, rights, interrupted);
231size_t UnixSocketConnection::Stream::writeVectors(
const struct iovec* vectors,
size_t vectorCount,
236 *interrupted =
false;
239 size_t firstVector = 0;
240 while (firstVector < vectorCount && !vectors[firstVector].iov_len) {
243 if (firstVector == vectorCount) {
247 Control* pendingControl = rights ?
new Control(0, rights) : nullptr;
248 StreamSerializationGuard sendGuard(m_SendLock, block ? StreamSerializationWait::Interruptible
249 : StreamSerializationWait::Nonblocking);
251 captureStreamInterruption(block, interrupted);
252 delete pendingControl;
256 size_t firstOffset = 0;
258 if (pendingControl) {
259 if (!m_Bytes.canWrite(block)) {
260 captureStreamInterruption(block, interrupted);
261 delete pendingControl;
265 ControlGuard controlGuard(*
this);
266 const uint8_t* first =
reinterpret_cast<const uint8_t*
>(vectors[firstVector].iov_base);
267 if (m_Bytes.write(first, 1, block) != 1) {
268 captureStreamInterruption(block, interrupted);
269 delete pendingControl;
273 pendingControl->byteOffset = m_BytesWritten;
274 m_Controls.pushBack(pendingControl);
280 for (
size_t i = firstVector; i < vectorCount; ++i) {
281 const uint8_t* buffer =
reinterpret_cast<const uint8_t*
>(vectors[i].iov_base);
282 const size_t count = vectors[i].iov_len;
283 const size_t offset = i == firstVector ? firstOffset : 0;
284 if (offset >= count) {
288 const size_t tail = m_Bytes.write(buffer + offset, count - offset, block);
290 ControlGuard controlGuard(*
this);
291 m_BytesWritten += tail;
294 if (tail < count - offset) {
295 captureStreamInterruption(block, interrupted);
303size_t UnixSocketConnection::Stream::read(uint8_t* buffer,
size_t count,
bool block,
305 struct iovec vector = {buffer, count};
306 return readVectors(&vector, 1, block, rights, interrupted);
309size_t UnixSocketConnection::Stream::readVectors(
struct iovec* vectors,
size_t vectorCount,
313 *interrupted =
false;
319 StreamSerializationGuard receiveGuard(m_ReceiveLock, block
320 ? StreamSerializationWait::Interruptible
321 : StreamSerializationWait::Nonblocking);
323 captureStreamInterruption(block, interrupted);
326 size_t totalRead = 0;
327 bool canBlock = block;
328 for (
size_t i = 0; i < vectorCount; ++i) {
329 uint8_t* buffer =
reinterpret_cast<uint8_t*
>(vectors[i].iov_base);
330 size_t count = vectors[i].iov_len;
335 if (!m_Bytes.canRead(canBlock)) {
336 captureStreamInterruption(canBlock, interrupted);
340 ControlGuard controlGuard(*
this);
341 while (m_Controls.count()) {
342 Control* stale = *m_Controls.begin();
343 if (stale->byteOffset >= m_BytesRead) {
346 delete m_Controls.popFront();
349 size_t amount = count;
350 if (rights && m_Controls.count()) {
351 Control* next = *m_Controls.begin();
352 const uint64_t distance = next->byteOffset - m_BytesRead;
353 if (distance < amount) {
354 amount =
static_cast<size_t>(distance + 1);
358 const size_t bytesRead = m_Bytes.read(buffer, amount, canBlock);
359 if (bytesRead < amount) {
360 captureStreamInterruption(canBlock, interrupted);
362 m_BytesRead += bytesRead;
363 bool consumedControl =
false;
364 while (bytesRead && m_Controls.count()) {
365 Control* crossed = *m_Controls.begin();
366 if (crossed->byteOffset >= m_BytesRead) {
369 crossed = m_Controls.popFront();
371 *rights = crossed->rights;
372 consumedControl =
true;
380 totalRead += bytesRead;
382 if (consumedControl || bytesRead < amount) {
390bool UnixSocketConnection::Stream::writePacket(
const uint8_t* buffer,
size_t count,
bool block,
394 *interrupted =
false;
396 StreamSerializationGuard sendGuard(m_SendLock, block ? StreamSerializationWait::Interruptible
397 : StreamSerializationWait::Nonblocking);
398 if (!sendGuard || !m_Records.canWrite(block)) {
399 captureStreamInterruption(block, interrupted);
403 Packet* packet =
new Packet();
404 packet->length = count;
405 packet->rights = rights;
407 packet->bytes =
new uint8_t[count];
408 MemoryCopy(packet->bytes, buffer, count);
415 ControlGuard controlGuard(*
this);
416 const uintptr_t record =
reinterpret_cast<uintptr_t
>(packet);
417 if (m_Records.writeAtomic(
reinterpret_cast<const uint8_t*
>(&record),
sizeof(record),
true) !=
420 captureStreamInterruption(block, interrupted);
423 m_PendingPackets.pushBack(packet);
427bool UnixSocketConnection::Stream::readPacket(uint8_t* buffer,
size_t count,
bool block,
429 uint64_t& bytesRead, uint64_t& packetLength,
432 bytesRead = packetLength = 0;
434 *interrupted =
false;
436 StreamSerializationGuard receiveGuard(m_ReceiveLock, block
437 ? StreamSerializationWait::Interruptible
438 : StreamSerializationWait::Nonblocking);
439 if (!receiveGuard || !m_Records.canRead(block)) {
440 captureStreamInterruption(block, interrupted);
444 ControlGuard controlGuard(*
this);
445 uintptr_t record = 0;
446 if (m_Records.read(
reinterpret_cast<uint8_t*
>(&record),
sizeof(record),
false) !=
448 captureStreamInterruption(block, interrupted);
451 Packet* packet =
reinterpret_cast<Packet*
>(record);
452 assert(m_PendingPackets.count() && *m_PendingPackets.begin() == packet);
453 m_PendingPackets.popFront();
454 packetLength = packet->length;
455 bytesRead = count < packetLength ? count : packetLength;
457 MemoryCopy(buffer, packet->bytes, bytesRead);
459 rights = packet->rights;
464bool UnixSocketConnection::Stream::canWrite(
bool block) {
465 return m_Packets ? m_Records.canWrite(block) : m_Bytes.canWrite(block);
468bool UnixSocketConnection::Stream::canRead(
bool block) {
469 return m_Packets ? m_Records.canRead(block) : m_Bytes.canRead(block);
472uint64_t UnixSocketConnection::Stream::readableGeneration()
const {
473 return m_Packets ? m_Records.readableGeneration() : m_Bytes.readableGeneration();
476uint64_t UnixSocketConnection::Stream::writableGeneration()
const {
477 return m_Packets ? m_Records.writableGeneration() : m_Bytes.writableGeneration();
480void UnixSocketConnection::Stream::disableWrites() {
481 m_Bytes.disableWrites();
482 m_Records.disableWrites();
485void UnixSocketConnection::Stream::disableReads() {
491 m_Bytes.disableReads();
492 m_Records.disableReads();
494#if !defined(PEDIGREE_EXTERNAL_SOURCE)
495 if (m_ControlLock.isOwnedByCurrentThread()) {
496 m_DiscardControlsRequested =
true;
501 ControlGuard controlGuard(*
this);
503 m_DiscardControlsRequested =
false;
508 m_Records.monitor(
waiter);
514void UnixSocketConnection::Stream::monitor(
Thread* thread,
Event* event) {
516 m_Records.monitor(thread, event);
518 m_Bytes.monitor(thread, event);
522void UnixSocketConnection::Stream::cullMonitorTargets(
Semaphore*
waiter) {
524 m_Records.cullMonitorTargets(
waiter);
526 m_Bytes.cullMonitorTargets(
waiter);
530void UnixSocketConnection::Stream::cullMonitorTargets(
Event* event) {
532 m_Records.cullMonitorTargets(event);
534 m_Bytes.cullMonitorTargets(event);
538void UnixSocketConnection::Stream::discardControls() {
539 while (m_PendingPackets.count()) {
540 delete m_PendingPackets.popFront();
542 while (m_Controls.count()) {
543 delete m_Controls.popFront();
547void UnixSocketConnection::Stream::discardControlsIfRequested() {
548 while (m_DiscardControlsRequested.compareAndSwap(
true,
false)) {
553UnixSocketConnection::UnixSocketConnection(
bool packets)
554 : m_FirstStream(packets),
555 m_SecondStream(packets),
558 m_Closed{false, false},
559 m_ReadShutdown{false, false},
560 m_WriteShutdown{false, false},
562 for (
size_t i = 0; i < 2; ++i) {
571 :
File(name, 0, 0, 0, 0, pFs, 0, pParent),
574 m_Datagrams(MAX_UNIX_DGRAM_BACKLOG),
575 m_Stream(MAX_UNIX_STREAM_QUEUE),
577 m_ConnectionSide(false),
582 if (m_Type == Datagram) {
592UnixSocket::~UnixSocket() {
597 if (m_Type != Datagram) {
600 bool shutdown =
false;
603 state = getStateLocked();
604 connection = m_Connection;
606 const bool side = m_ConnectionSide;
609 ? connection->m_WriteShutdown[side] || connection->m_ReadShutdown[side ? 0 : 1]
610 : connection->m_ReadShutdown[side] || connection->m_WriteShutdown[side ? 0 : 1];
614 if (state == Listening) {
615 return !bWriting && m_Stream.
canRead(timeout == 1);
618 if (state == Closed) {
622 if (state != Active || !connection) {
631 if (outgoingStream(connection)->canWrite(timeout == 1)) {
635 if (incomingStream(connection)->canRead(timeout == 1)) {
644 if (m_State == Closed) {
650 return m_Datagrams.
waitFor(bWriting ? RingBufferWait::Writing : RingBufferWait::Reading);
651 }
else if (bWriting) {
662 return recvfrom(size, buffer, bCanBlock, remote);
665uint64_t UnixSocket::recvfrom(uint64_t size, uintptr_t buffer,
bool bCanBlock,
String& from) {
666 if (m_Type == SequencedPacket) {
669 uint64_t bytesRead = 0, packetLength = 0;
670 receivePacket(size, buffer, bCanBlock, rights, bytesRead, packetLength);
673 if (m_Type == Streaming) {
679 uint64_t bytesRead = 0;
680 uint64_t datagramLength = 0;
681 receiveDatagram(size, buffer, bCanBlock, from, rights, bytesRead, datagramLength);
687 struct iovec vector = {
reinterpret_cast<void*
>(buffer),
688 static_cast<size_t>(size > SIZE_MAX ? SIZE_MAX : size)};
689 return receiveStream(&vector, 1, bCanBlock, rights, interrupted);
695 *interrupted =
false;
705 state = getStateLocked();
706 connection = m_Connection;
709 if (m_Type != Streaming || !connection || (state != Active && state != Closed)) {
713 return incomingStream(connection)
714 ->readVectors(vectors, vectorCount, state == Active && bCanBlock, rights, interrupted);
719 uint64_t& datagramLength) {
726 if (m_State == Closed || m_Type != Datagram) {
735 }
else if (!
select(
false, 0)) {
739 struct buf* datagram =
nullptr;
740 DatagramBuffer::Error error = DatagramBuffer::NoError;
741 if (!m_Datagrams.
read(datagram, error)) {
745 datagramLength = datagram->len;
746 bytesRead = size < datagramLength ? size : datagramLength;
748 MemoryCopy(
reinterpret_cast<void*
>(buffer), datagram->pBuffer, bytesRead);
750 if (datagram->remotePath) {
751 from.assign(datagram->remotePath, datagram->remotePathLen);
753 rights = datagram->rights;
754 destroyDatagram(datagram);
760 if (m_Type == SequencedPacket) {
762 return sendPacket(size, buffer, bCanBlock, rights) ? size : 0;
764 if (m_Type == Streaming) {
766 return sendStream(size, buffer, bCanBlock, rights);
770 return sendDatagram(size, buffer, bCanBlock, location, rights) ? size : 0;
775 struct iovec vector = {
reinterpret_cast<void*
>(buffer),
776 static_cast<size_t>(size > SIZE_MAX ? SIZE_MAX : size)};
777 return sendStream(&vector, 1, bCanBlock, rights, interrupted);
783 *interrupted =
false;
789 state = getStateLocked();
790 connection = m_Connection;
793 if (m_Type != Streaming || !connection || state != Active) {
794 N_NOTICE(
"UnixSocket::write => closed or not connected");
798 return outgoingStream(connection)
799 ->writeVectors(vectors, vectorCount, bCanBlock, rights, interrupted);
807 auto fail = [error](
int value) {
813 if (m_Type != SequencedPacket) {
814 return fail(EPROTOTYPE);
816 if (size > MAX_UNIX_STREAM_QUEUE) {
817 return fail(EMSGSIZE);
823 state = getStateLocked();
824 connection = m_Connection;
826 if (!connection || state != Active || writeShutdown()) {
827 return fail(wasConnected() ? EPIPE : ENOTCONN);
829 bool interrupted =
false;
830 if (outgoingStream(connection)
831 ->writePacket(
reinterpret_cast<const uint8_t*
>(buffer), size, bCanBlock, rights,
838 return fail(getState() == Closed || writeShutdown() ? EPIPE : EAGAIN);
841bool UnixSocket::receivePacket(uint64_t size, uintptr_t buffer,
bool bCanBlock,
843 uint64_t& packetLength,
bool* interrupted) {
845 bytesRead = packetLength = 0;
847 *interrupted =
false;
853 state = getStateLocked();
854 connection = m_Connection;
856 if (m_Type != SequencedPacket || !connection || (state != Active && state != Closed)) {
859 return incomingStream(connection)
860 ->readPacket(
reinterpret_cast<uint8_t*
>(buffer), size, state == Active && bCanBlock, rights,
861 bytesRead, packetLength, interrupted);
869 auto fail = [error](
int value) {
875 if (m_Type != Datagram) {
876 return fail(EPROTOTYPE);
878 if (getState() == Closed) {
879 return fail(ECONNREFUSED);
882 struct buf* b =
new buf();
887 b->pBuffer =
new char[size];
892 MemoryCopy(b->pBuffer,
reinterpret_cast<void*
>(buffer), size);
897 b->remotePathLen = StringLength(
reinterpret_cast<const char*
>(source));
898 b->remotePath =
new char[b->remotePathLen + 1];
899 if (!b->remotePath) {
903 MemoryCopy(b->remotePath,
reinterpret_cast<const void*
>(source), b->remotePathLen + 1);
906 const DatagramBuffer::Error result = bCanBlock ? m_Datagrams.
write(b) : m_Datagrams.
tryWrite(b);
907 if (result != DatagramBuffer::NoError) {
909 if (result == DatagramBuffer::Closed) {
910 return fail(ECONNREFUSED);
912 if (result == DatagramBuffer::Interrupted || result == DatagramBuffer::ThreadTerminating) {
922void UnixSocket::destroyDatagram(
struct buf* datagram) {
927 delete[] datagram->remotePath;
928 delete[] datagram->pBuffer;
932bool UnixSocket::bind(
UnixSocket* other,
bool block) {
935 if (!other || m_Type == Datagram || other->m_Type != m_Type) {
943 if (m_State != Inactive || other->m_State != Inactive || m_Connection || other->m_Connection) {
944 N_NOTICE(
"UnixSocket::bind endpoints are not inactive");
948 m_Connection = connection;
949 m_ConnectionSide =
false;
950 other->m_Connection = connection;
951 other->m_ConnectionSide =
true;
952 m_State = Connecting;
953 other->m_State = Connecting;
956 connection->m_Creds[0] = m_Creds;
962void UnixSocket::unbind() {
968 connection = m_Connection;
970 side = m_ConnectionSide;
971 connection->m_Closed[side] =
true;
975 while (m_PendingSockets.
count()) {
980 N_NOTICE(
"UnixSocket::unbind");
982 if (m_Type == Datagram) {
984 struct buf* datagram =
nullptr;
986 destroyDatagram(datagram);
992 side ? &connection->m_SecondStream : &connection->m_FirstStream;
994 side ? &connection->m_FirstStream : &connection->m_SecondStream;
995 incoming->disableReads();
996 outgoing->disableWrites();
997 notifyStream(incoming);
998 notifyStream(outgoing);
1005 while (pending.
count()) {
1007 socket->failConnection();
1012bool UnixSocket::shutdown(
int how) {
1017 if (m_Type == Datagram || !m_Connection || !m_Connection->m_Active || m_Connection->m_Failed ||
1018 m_Connection->m_Closed[0] || m_Connection->m_Closed[1]) {
1019 SYSCALL_ERROR(NotConnected);
1023 connection = m_Connection;
1024 side = m_ConnectionSide;
1025 if (how == SHUT_RD || how == SHUT_RDWR) {
1026 connection->m_ReadShutdown[side] =
true;
1028 if (how == SHUT_WR || how == SHUT_RDWR) {
1029 connection->m_WriteShutdown[side] =
true;
1033 auto* incoming = side ? &connection->m_SecondStream : &connection->m_FirstStream;
1034 auto* outgoing = side ? &connection->m_FirstStream : &connection->m_SecondStream;
1035 if (how == SHUT_RD || how == SHUT_RDWR) {
1036 incoming->disableReads();
1037 notifyStream(incoming);
1039 if (how == SHUT_WR || how == SHUT_RDWR) {
1040 outgoing->disableWrites();
1041 notifyStream(outgoing);
1046bool UnixSocket::writeShutdown()
const {
1048 if (!m_Connection) {
1051 const bool side = m_ConnectionSide;
1052 return m_Connection->m_WriteShutdown[side] || m_Connection->m_ReadShutdown[side ? 0 : 1];
1055bool UnixSocket::readShutdown()
const {
1057 if (!m_Connection) {
1060 const bool side = m_ConnectionSide;
1061 return m_Connection->m_ReadShutdown[side] || m_Connection->m_WriteShutdown[side ? 0 : 1];
1064void UnixSocket::acknowledgeBind() {
1068 connection = m_Connection;
1069 if (!connection || connection->m_Failed || connection->m_Closed[0] || connection->m_Closed[1] ||
1070 connection->m_Active) {
1074 N_NOTICE(
"acking bind");
1076 connection->m_Active =
true;
1080 connection->m_Creds[m_ConnectionSide ? 1 : 0] = m_Creds;
1083 notifyStream(&connection->m_FirstStream);
1084 notifyStream(&connection->m_SecondStream);
1087bool UnixSocket::addSocket(
UnixSocket* socket) {
1090 if (m_State != Listening || !socket || !socket->m_Connection || socket->m_Connection->m_Failed ||
1091 socket->m_Connection->m_Closed[0] || socket->m_Connection->m_Closed[1]) {
1095 connection = socket->m_Connection;
1096 socket->m_Creds = m_Creds;
1097 connection->m_Creds[socket->m_ConnectionSide ? 1 : 0] = m_Creds;
1098 connection->m_Active =
true;
1099 socket->m_State = Active;
1102 N_NOTICE(
"adding listening socket");
1108 if (m_Stream.
write(&c, 1,
false) == 1) {
1109 notifyStream(&connection->m_FirstStream);
1110 notifyStream(&connection->m_SecondStream);
1116 if (*it == socket) {
1117 m_PendingSockets.
erase(it);
1124UnixSocket* UnixSocket::getSocket(
bool block) {
1126 if (m_Stream.
read(&c, 1, block) != 1) {
1130 N_NOTICE(
"got a socket");
1134 if (m_State != Listening || !m_PendingSockets.
count()) {
1138 N_NOTICE(
"popping socket");
1139 return m_PendingSockets.
popFront();
1143 if (m_Type == Datagram) {
1150 bool closed =
false;
1153 closed = getStateLocked() == Closed;
1154 connection = m_Connection;
1155 side = m_ConnectionSide;
1167 if (monitorRead ||
write) {
1172 closed = getStateLocked() == Closed;
1181 side ? &connection->m_SecondStream : &connection->m_FirstStream;
1183 side ? &connection->m_FirstStream : &connection->m_SecondStream;
1185 incoming->monitor(
waiter);
1187 if (
write && (!monitorRead || outgoing != incoming)) {
1188 outgoing->monitor(
waiter);
1193 closed = getStateLocked() == Closed;
1199 notifyStream(incoming);
1201 if (
write && (!monitorRead || outgoing != incoming)) {
1202 notifyStream(outgoing);
1208 if (m_Type == Datagram) {
1216 connection = m_Connection;
1223 connection->m_FirstStream.cullMonitorTargets(
waiter);
1224 connection->m_SecondStream.cullMonitorTargets(
waiter);
1227void UnixSocket::addWaiter(
Thread* thread,
Event* event) {
1228 if (m_Type == Datagram) {
1229 m_Datagrams.
monitor(thread, event);
1234 bool closed =
false;
1237 closed = getStateLocked() == Closed;
1238 connection = m_Connection;
1244#if !defined(PEDIGREE_EXTERNAL_SOURCE)
1251 m_Stream.
monitor(thread, event);
1254 closed = getStateLocked() == Closed;
1264 first->monitor(thread, event);
1265 if (second != first) {
1266 second->monitor(thread, event);
1271 closed = getStateLocked() == Closed;
1276 notifyStream(first);
1277 if (second != first) {
1278 notifyStream(second);
1283void UnixSocket::removeWaiter(
Event* event) {
1284 if (m_Type == Datagram) {
1292 connection = m_Connection;
1299 connection->m_FirstStream.cullMonitorTargets(event);
1300 connection->m_SecondStream.cullMonitorTargets(event);
1303bool UnixSocket::markListening() {
1306 if (m_Type == Datagram) {
1311 if (m_State != Inactive) {
1317 m_State = Listening;
1321UnixSocket::SocketState UnixSocket::getState()
const {
1323 return getStateLocked();
1326UnixSocket::SocketState UnixSocket::getStateLocked()
const {
1327 if (m_Type == Datagram || !m_Connection) {
1331 if (m_Connection->m_Failed || m_Connection->m_Closed[m_ConnectionSide ? 1 : 0] ||
1332 m_Connection->m_Closed[m_ConnectionSide ? 0 : 1]) {
1336 return m_Connection->m_Active ? Active : Connecting;
1339bool UnixSocket::wasConnected()
const {
1341 return m_Connection && m_Connection->m_Active && !m_Connection->m_Failed;
1346 if (m_Type == Datagram) {
1347 generations.read = m_Datagrams.readableGeneration();
1356 state = getStateLocked();
1357 connection = m_Connection;
1358 side = m_ConnectionSide;
1361 if (state == Listening || !connection) {
1363 generations.write = m_Stream.writableGeneration();
1368 side ? &connection->m_SecondStream : &connection->m_FirstStream;
1370 side ? &connection->m_FirstStream : &connection->m_SecondStream;
1371 generations.read = incoming->readableGeneration();
1372 generations.write = outgoing->writableGeneration();
1376void UnixSocket::failConnection() {
1380 connection = m_Connection;
1386 connection->m_Failed =
true;
1387 connection->m_Closed[0] =
true;
1388 connection->m_Closed[1] =
true;
1392 connection->m_FirstStream.disableWrites();
1393 connection->m_FirstStream.disableReads();
1394 connection->m_SecondStream.disableWrites();
1395 connection->m_SecondStream.disableReads();
1396 notifyStream(&connection->m_FirstStream);
1397 notifyStream(&connection->m_SecondStream);
1400struct ucred
UnixSocket::getPeerCredentials() const {
1402 if (!m_Connection) {
1410 return m_Connection->m_Creds[m_ConnectionSide ? 0 : 1];
1415 return m_ConnectionSide ? &connection->m_SecondStream : &connection->m_FirstStream;
1420 return m_ConnectionSide ? &connection->m_FirstStream : &connection->m_SecondStream;
1423void UnixSocket::setCreds() {
1426 m_Creds.uid = pCurrentProcess->
getUserId();
1427 m_Creds.gid = pCurrentProcess->getGroupId();
1433 :
Directory(name, 0, 0, 0, 0, pFs, 0, pParent), m_Lock() {
1434 cacheDirectoryContents();
1437UnixDirectory::~UnixDirectory() {}
1439bool UnixDirectory::addEntry(
const String& filename,
File* pFile) {
1443bool UnixDirectory::removeEntry(
const String& filename,
File* pFile) {
1449 if (parent ==
this) {
1450 SYSCALL_ERROR(InvalidArgument);
1457 SYSCALL_ERROR(IoError);
1461 SYSCALL_ERROR(NotEmpty);
1464 if (!parent->removeEntry(filename,
this))
1474UnixFilesystem::UnixFilesystem() :
Filesystem(), m_pRoot(0) {
1481 m_pRoot->setPermissions(FILE_UR | FILE_UW | FILE_UX | FILE_GR | FILE_GW | FILE_GX | FILE_OR |
1485UnixFilesystem::~UnixFilesystem() {
1488 ERROR(
"UnixFilesystem::~UnixFilesystem: root didn't get destroyed");
1492Mutex& UnixFilesystem::namespaceLock() {
1493 return m_NamespaceLock;
1500 if (!pParent->addEntry(filename, pSocket)) {
1506 pSocket->setPermissions(FILE_UR | FILE_UW | FILE_UX | FILE_GR | FILE_GW | FILE_GX | FILE_OR |
1516 if (!pParent->addEntry(filename, pChild)) {
1522 pChild->setPermissions(FILE_UR | FILE_UW | FILE_UX | FILE_GR | FILE_GW | FILE_GX | FILE_OR |
1531 return static_cast<UnixDirectory*
>(file)->removeFromParent(pParent, filename);
1533 return pParent->removeEntry(filename, file);
1537 if (m_Type == SequencedPacket) {
size_t write(const T *buffer, size_t count, bool block=true)
void monitor(Thread *pThread, Event *pEvent)
void cullMonitorTargets(Thread *pThread)
uint64_t readableGeneration() const
size_t read(T *buffer, size_t count, bool block=true)
static Directory * fromFile(File *pF)
bool addDirectoryEntry(const String &name, File *pTarget)
void markCachePopulated()
Mutex & namespaceMutationLock()
bool removeDirectoryEntry(const HashedStringView &name, File *expected)
ReadStatus isEmpty(bool &empty)
virtual uint64_t read(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock=true) final
virtual uint64_t write(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock=true) final
virtual bool isDirectory()
::Iterator< T, node_t > Iterator
MUST_USE_RESULT bool waitFor(MailboxWait::WaitType wait, Time::Timestamp &timeout, Error &error)
waitFor - block until the given condition is true (readable/writeable)
bool canWrite()
canWrite - is it possible to write to the ring buffer without blocking?
MUST_USE_RESULT bool takeAfterClose(T &out)
MUST_USE_RESULT bool read(T &out, Time::Timestamp &timeout, Error &error)
Error tryWrite(const T &obj)
Error write(const T &obj, Time::Timestamp &timeout)
Writes one object, waiting for space up to the supplied timeout.
bool dataReady()
dataReady - is data ready for reading from the ring buffer?
size_t getUserspaceId() const
virtual int64_t getUserId() const
static ProcessorInformation & information()
bool monitor(Thread *pThread, Event *pEvent)
monitor - add a new Event to be fired when something happens
void cullMonitorTargets(Thread *pThread)
Cull all monitor targets pointing to pThread.
static bool create(size_t descriptorCount, SharedPointer< SocketRights > &rights)
bool sendEvent(Event *pEvent)
virtual void cacheDirectoryContents()
virtual bool createFile(File *parent, const String &filename, uint32_t mask)
virtual bool removeNode(File *parent, const String &filename, File *file)
virtual bool createDirectory(File *parent, const String &filename, uint32_t mask)
uint64_t sendStream(uint64_t size, uintptr_t buffer, bool bCanBlock, const SharedPointer< SocketRights > &rights, bool *interrupted=nullptr)
bool sendDatagram(uint64_t size, uintptr_t buffer, bool bCanBlock, uintptr_t source, const SharedPointer< SocketRights > &rights, int *error=nullptr)
virtual uint64_t writeBytewise(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock=true)
virtual uint64_t readBytewise(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock=true)
uint64_t receiveStream(uint64_t size, uintptr_t buffer, bool bCanBlock, SharedPointer< SocketRights > *rights, bool *interrupted=nullptr)
virtual int select(bool bWriting=false, int timeout=0)
bool receiveDatagram(uint64_t size, uintptr_t buffer, bool bCanBlock, String &from, SharedPointer< SocketRights > &rights, uint64_t &bytesRead, uint64_t &datagramLength)
ReadinessGenerations readinessGenerations() override
bool sendPacket(uint64_t size, uintptr_t buffer, bool bCanBlock, const SharedPointer< SocketRights > &rights, int *error=nullptr)
void trackFile(File *pFile)
Track a File object that exists. It is necessary to keep track of File objects, or at least those tha...
Iterator erase(Iterator &Iter)
void pushBack(const T &value)
void pushBack(const T &value)
void clear(bool freeMem=false)