The Pedigree Project 0.1
pipe-transfer-syscalls.cc
1/* Copyright (c) 2026, Pedigree Developers. */
2#include "pedigree/kernel/process/Process.h"
3#include "pedigree/kernel/process/TerminationDeferral.h"
4#include "pedigree/kernel/process/Thread.h"
5#include "pedigree/kernel/processor/Processor.h"
6#include "pedigree/kernel/processor/ProcessorInformation.h"
7#include "pedigree/kernel/syscallError.h"
8#include "pedigree/kernel/utilities/Pointers.h"
9#include "pedigree/kernel/utilities/assert.h"
10
11#include <fcntl.h>
12
13#include "FileDescriptor.h"
14#include "PosixSubsystem.h"
15#include "modules/system/vfs/MountView.h"
16#include "modules/system/vfs/Pipe.h"
17#include "net-syscalls.h"
18#include "pipe-transfer-syscalls.h"
19
20namespace {
21constexpr unsigned KnownFlags = 0xf;
22constexpr unsigned Nonblock = 2;
23constexpr uint64_t MaximumPosition = 0x7fffffffffffffffULL;
24using Status = PipeBuffer::Status;
25using PositionGuard = FileDescriptor::TransferPositionGuard;
26using PositionSide = PositionGuard::Endpoint;
27
28struct Endpoint {
29 DescriptorLease descriptor;
32 File* file = nullptr;
33 Pipe* pipe = nullptr;
34 int flags = 0;
35};
36
37bool acquireEndpoint(PosixSubsystem* subsystem, int fd, Endpoint& endpoint) {
38 if (!subsystem || !subsystem->acquireFileDescriptor(fd, endpoint.descriptor)) {
39 SYSCALL_ERROR(BadFileDescriptor);
40 return false;
41 }
42 endpoint.description = endpoint.descriptor->acquireOpenFileDescription();
43 if (!endpoint.description) {
44 SYSCALL_ERROR(BadFileDescriptor);
45 return false;
46 }
47 endpoint.flags = endpoint.descriptor->getStatusFlags();
48 if (endpoint.flags & O_PATH) {
49 SYSCALL_ERROR(BadFileDescriptor);
50 return false;
51 }
52 endpoint.file = endpoint.description->getFile();
53 endpoint.socket = endpoint.description->getNetworkImpl();
54 if (endpoint.file && (endpoint.file->isPipe() || endpoint.file->isFifo()))
55 endpoint.pipe = Pipe::fromFile(endpoint.file);
56 return true;
57}
58
59bool accessAllowed(const Endpoint& endpoint, bool writing) {
60 const int mode = endpoint.flags & O_ACCMODE;
61 if (endpoint.socket || mode == O_RDWR || mode == (writing ? O_WRONLY : O_RDONLY))
62 return true;
63 SYSCALL_ERROR(BadFileDescriptor);
64 return false;
65}
66
67bool validRange(uint64_t offset, size_t count) {
68 if (offset > MaximumPosition || count > MaximumPosition - offset) {
69 SYSCALL_ERROR(InvalidArgument);
70 return false;
71 }
72 return true;
73}
74
75bool interrupted(Thread* thread) {
76 if (thread->getInterruptionReason() == Thread::InterruptedBySignal ||
77 thread->getUnwindState() != Thread::Continue) {
78 SYSCALL_ERROR(Interrupted);
79 return true;
80 }
81 return false;
82}
83
84ssize_t pipeResult(PipeBuffer::Result result, bool writing, bool& pipeSignal) {
85 if (result.count)
86 return static_cast<ssize_t>(result.count);
87 switch (result.status) {
88 case Status::Ready:
89 case Status::Eof:
90 return 0;
91 case Status::Closed:
92 if (!writing)
93 return 0;
94 pipeSignal = true;
95 SYSCALL_ERROR(BrokenPipe);
96 break;
97 case Status::WouldBlock:
98 SYSCALL_ERROR(NoMoreProcesses);
99 break;
100 case Status::Interrupted:
101 SYSCALL_ERROR(Interrupted);
102 break;
103 case Status::Invalid:
104 SYSCALL_ERROR(InvalidArgument);
105 break;
106 }
107 return -1;
108}
109
110ssize_t finish(Thread* thread, PosixSubsystem* subsystem, ssize_t result, bool pipeSignal) {
111 const size_t error = result < 0 ? thread->getErrno() : 0;
112 thread->clearInterruption();
113 if (pipeSignal)
114 subsystem->threadException(thread, Subsystem::Pipe);
115 thread->setErrno(error);
116 return result;
117}
118
119ssize_t mixedSplice(Thread* thread, Endpoint& input, Endpoint& output, int64_t* explicitPosition,
120 size_t count, unsigned flags, bool& pipeSignal) {
121 const bool writingPipe = output.pipe != nullptr;
122 Pipe* pipe = writingPipe ? output.pipe : input.pipe;
123 Endpoint& regular = writingPipe ? input : output;
124 const bool socket = !writingPipe && bool(output.socket);
125 if (socket) {
126 if (explicitPosition || output.socket->getType() != SOCK_STREAM ||
127 (output.socket->getDomain() != AF_UNIX && output.socket->getDomain() != AF_INET &&
128 output.socket->getDomain() != AF_INET6)) {
129 SYSCALL_ERROR(InvalidArgument);
130 return -1;
131 }
132 } else if (!regular.file || !regular.file->supportsRegularFileOperations()) {
133 SYSCALL_ERROR(InvalidArgument);
134 return -1;
135 }
136 if (!writingPipe && (output.flags & O_APPEND)) {
137 SYSCALL_ERROR(InvalidArgument);
138 return -1;
139 }
140 const uint64_t initial = explicitPosition ? static_cast<uint64_t>(*explicitPosition)
141 : socket ? 0
142 : regular.descriptor->getOffset();
143 if (!validRange(initial, count))
144 return -1;
145 if (socket) {
146 struct sockaddr_storage peer = {};
147 socklen_t length = sizeof(peer);
148 thread->setErrno(0);
149 if (output.socket->getpeername(&peer, &length) < 0)
150 return -1;
151 }
152 const size_t maximum = count < PipeBuffer::Capacity ? count : PipeBuffer::Capacity;
153 auto scratch = UniqueArray<uint8_t>::allocate(maximum);
154 if (!scratch) {
155 SYSCALL_ERROR(OutOfMemory);
156 return -1;
157 }
158 const bool canBlock =
159 !(flags & Nonblock) && !((writingPipe ? output.flags : input.flags) & O_NONBLOCK);
160
161 for (;;) {
162 if (interrupted(thread))
163 return -1;
164 // No position or writer lock may keep the opposite transfer from draining.
165 const auto ready = pipe->waitTransfer(writingPipe, canBlock);
166 if (ready.status != Status::Ready)
167 return pipeResult(ready, writingPipe, pipeSignal);
168 bool retry = false;
169 const ssize_t result = [&]() -> ssize_t {
170 PositionGuard positions(input.description, output.description,
171 writingPipe && !explicitPosition, !writingPipe && !socket);
172 if (!writingPipe && (positions.statusFlags(PositionSide::Output) & O_APPEND)) {
173 SYSCALL_ERROR(InvalidArgument);
174 return -1;
175 }
176 uint64_t position =
177 explicitPosition ? static_cast<uint64_t>(*explicitPosition)
178 : socket ? 0
179 : positions.offset(writingPipe ? PositionSide::Input : PositionSide::Output);
180 if (!validRange(position, count) || interrupted(thread))
181 return -1;
182 VfsMountView::WriteLease mountWrite;
183 if (output.file && !output.pipe && output.descriptor->openingPath() &&
184 !mountWrite.acquire(output.descriptor->openingPath())) {
185 return -1;
186 }
187 ssize_t moved;
188 if (writingPipe) {
189 Pipe::WriteReservation reservation;
190 const auto reserved = pipe->reserveWrite(maximum, false, reservation);
191 if (reserved.status == Status::WouldBlock) {
192 retry = true;
193 return -1;
194 }
195 if (reserved.status != Status::Ready)
196 return pipeResult(reserved, true, pipeSignal);
197 thread->setErrno(0);
198 const size_t read = regular.file->read(position, reservation.size(),
199 reinterpret_cast<uintptr_t>(scratch.get()), true);
200 assert(read <= reservation.size());
201 if (!read)
202 return thread->getErrno() || interrupted(thread) ? -1 : 0;
203 if (interrupted(thread))
204 return -1;
205 moved = pipeResult(reservation.commit(scratch.get(), read), true, pipeSignal);
206 } else {
207 auto fromPipe = [&](File::WriteGuard* writer) -> ssize_t {
208 Pipe::ReadReservation reservation;
209 const auto reserved = pipe->reserveRead(maximum, false, reservation);
210 if (reserved.status == Status::WouldBlock) {
211 retry = true;
212 return -1;
213 }
214 if (reserved.status != Status::Ready)
215 return pipeResult(reserved, false, pipeSignal);
216 size_t amount = reservation.size();
217 if (writer) {
218 const uint64_t limit = regular.file->maximumFileSize();
219 if (position >= limit) {
220 SYSCALL_ERROR(FileTooLarge);
221 return -1;
222 }
223 if (amount > limit - position)
224 amount = static_cast<size_t>(limit - position);
225 }
226 reservation.copyTo(scratch.get(), amount);
227 if (interrupted(thread))
228 return -1;
229 thread->setErrno(0);
230 const ssize_t written =
231 writer ? static_cast<ssize_t>(writer->write(
232 position, amount, reinterpret_cast<uintptr_t>(scratch.get()), true))
233 : posix_send_descriptor(output.descriptor, scratch.get(), amount, 0, true);
234 if (!writer && thread->getErrno() == Error::BrokenPipe)
235 pipeSignal = true;
236 if (written <= 0) {
237 if (!thread->getErrno() && !interrupted(thread))
238 SYSCALL_ERROR(IoError);
239 return -1;
240 }
241 assert(static_cast<size_t>(written) <= amount);
242 reservation.consume(static_cast<size_t>(written));
243 return written;
244 };
245 if (socket) {
246 moved = fromPipe(nullptr);
247 } else {
248 auto writer = regular.file->lockWrites();
249 moved = fromPipe(&writer);
250 }
251 }
252 if (moved > 0) {
253 position += static_cast<size_t>(moved);
254 if (explicitPosition)
255 *explicitPosition = static_cast<int64_t>(position);
256 else if (!socket)
257 positions.commitOffset(writingPipe ? PositionSide::Input : PositionSide::Output,
258 position);
259 thread->setErrno(0);
260 }
261 return moved;
262 }();
263 if (!retry)
264 return result;
265 }
266}
267} // namespace
268
269ssize_t posix_splice(int inputFd, int64_t* inputOffset, int outputFd, int64_t* outputOffset,
270 size_t count, unsigned flags) {
271 TerminationDeferral lifetime;
272 Thread* thread = Processor::information().getCurrentThread();
273 auto* subsystem = static_cast<PosixSubsystem*>(thread->getParent()->getSubsystem());
274 thread->clearInterruption();
275 thread->setErrno(0);
276 if (!count)
277 return 0;
278 if (flags & ~KnownFlags) {
279 SYSCALL_ERROR(InvalidArgument);
280 return -1;
281 }
282 bool pipeSignal = false;
283 const ssize_t result = [&]() -> ssize_t {
284 Endpoint input, output;
285 if (!acquireEndpoint(subsystem, inputFd, input) ||
286 !acquireEndpoint(subsystem, outputFd, output))
287 return -1;
288 if ((input.pipe && inputOffset) || (output.pipe && outputOffset)) {
289 SYSCALL_ERROR(IllegalSeek);
290 return -1;
291 }
292 int64_t inputPosition = 0, outputPosition = 0;
293 if ((outputOffset &&
294 !PosixSubsystem::copyFromUser(&outputPosition, outputOffset, sizeof(outputPosition))) ||
295 (inputOffset &&
296 !PosixSubsystem::copyFromUser(&inputPosition, inputOffset, sizeof(inputPosition)))) {
297 SYSCALL_ERROR(BadAddress);
298 return -1;
299 }
300 if (!accessAllowed(input, false) || !accessAllowed(output, true))
301 return -1;
302 VfsMountView::WriteLease mountWrite;
303 if (output.file && !output.pipe && output.descriptor->openingPath() &&
304 !mountWrite.acquire(output.descriptor->openingPath())) {
305 return -1;
306 }
307 ssize_t moved;
308 if (input.pipe && output.pipe) {
309 const bool canBlock = !(flags & Nonblock) && !((input.flags | output.flags) & O_NONBLOCK);
310 moved =
311 pipeResult(input.pipe->transferTo(*output.pipe, count, true, canBlock), true, pipeSignal);
312 } else if (input.pipe || output.pipe) {
313 int64_t* position = input.pipe ? (outputOffset ? &outputPosition : nullptr)
314 : (inputOffset ? &inputPosition : nullptr);
315 moved = mixedSplice(thread, input, output, position, count, flags, pipeSignal);
316 } else {
317 SYSCALL_ERROR(InvalidArgument);
318 return -1;
319 }
320 if (moved < 0)
321 return moved;
322 // Unlike sendfile, splice exports offsets after EOF but never after an error.
323 if ((outputOffset &&
324 !PosixSubsystem::copyToUser(outputOffset, &outputPosition, sizeof(outputPosition))) ||
325 (inputOffset &&
326 !PosixSubsystem::copyToUser(inputOffset, &inputPosition, sizeof(inputPosition)))) {
327 SYSCALL_ERROR(BadAddress);
328 return -1;
329 }
330 return moved;
331 }();
332 return finish(thread, subsystem, result, pipeSignal);
333}
334
335ssize_t posix_tee(int inputFd, int outputFd, size_t count, unsigned flags) {
336 TerminationDeferral lifetime;
337 Thread* thread = Processor::information().getCurrentThread();
338 auto* subsystem = static_cast<PosixSubsystem*>(thread->getParent()->getSubsystem());
339 thread->clearInterruption();
340 thread->setErrno(0);
341 if (flags & ~KnownFlags) {
342 SYSCALL_ERROR(InvalidArgument);
343 return -1;
344 }
345 if (!count)
346 return 0;
347 bool pipeSignal = false;
348 const ssize_t result = [&]() -> ssize_t {
349 Endpoint input, output;
350 if (!acquireEndpoint(subsystem, inputFd, input) ||
351 !acquireEndpoint(subsystem, outputFd, output) || !accessAllowed(input, false) ||
352 !accessAllowed(output, true))
353 return -1;
354 if (!input.pipe || !output.pipe) {
355 SYSCALL_ERROR(InvalidArgument);
356 return -1;
357 }
358 const bool canBlock = !(flags & Nonblock) && !((input.flags | output.flags) & O_NONBLOCK);
359 return pipeResult(input.pipe->transferTo(*output.pipe, count, false, canBlock), true,
360 pipeSignal);
361 }();
362 return finish(thread, subsystem, result, pipeSignal);
363}
Definition File.h:75
Definition Pipe.h:36
static Pipe * fromFile(File *pF)
Definition Pipe.h:42
virtual void threadException(Thread *pThread, ExceptionType eType, InterruptState *pState=nullptr, uintptr_t faultAddress=0, uintptr_t errorCode=0)
bool acquireFileDescriptor(size_t fd, DescriptorLease &descriptor)
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()
void setErrno(size_t err)
Definition Thread.h:482
@ Continue
No unwind necessary, carry on as normal.
Definition Thread.h:517
size_t getErrno()
Definition Thread.h:477
UnwindType getUnwindState()
Definition Thread.h:535
Process * getParent() const
Definition Thread.h:340