From 3386b8b7858f5659c205b8f088eb568692699591 Mon Sep 17 00:00:00 2001 From: Augustin Cavalier Date: Tue, 25 Jul 2023 01:25:05 -0400 Subject: [PATCH] libbsd: Add a basic kqueue implementation. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit It only supports file descriptors and processes (threads), and a few flags (not all) to go with them. This has been tested extensively against libuv. Change-Id: I6fc5930fa7273698172c9c695965842b5df44f03 Reviewed-on: https://review.haiku-os.org/c/haiku/+/6746 Reviewed-by: waddlesplash Reviewed-by: Jérôme Duval --- headers/compatibility/bsd/sys/event.h | 118 ++++++++++++ src/libs/bsd/Jamfile | 3 + src/libs/bsd/kqueue.cpp | 263 ++++++++++++++++++++++++++ 3 files changed, 384 insertions(+) create mode 100644 headers/compatibility/bsd/sys/event.h create mode 100644 src/libs/bsd/kqueue.cpp diff --git a/headers/compatibility/bsd/sys/event.h b/headers/compatibility/bsd/sys/event.h new file mode 100644 index 0000000000..28bc13eb96 --- /dev/null +++ b/headers/compatibility/bsd/sys/event.h @@ -0,0 +1,118 @@ +/*- + * SPDX-License-Identifier: BSD-2-Clause + * + * Copyright (c) 1999,2000,2001 Jonathan Lemon + * All rights reserved. + * + * Redistribution and use in source and binary forms, with or without + * modification, are permitted provided that the following conditions + * are met: + * 1. Redistributions of source code must retain the above copyright + * notice, this list of conditions and the following disclaimer. + * 2. Redistributions in binary form must reproduce the above copyright + * notice, this list of conditions and the following disclaimer in the + * documentation and/or other materials provided with the distribution. + * + * THIS SOFTWARE IS PROVIDED BY THE AUTHOR AND CONTRIBUTORS ``AS IS'' AND + * ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE + * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE + * ARE DISCLAIMED. IN NO EVENT SHALL THE AUTHOR OR CONTRIBUTORS BE LIABLE + * FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL + * DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS + * OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) + * HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT + * LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY + * OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF + * SUCH DAMAGE. + * + * $FreeBSD$ + */ +#ifndef _BSD_SYS_EVENT_H_ +#define _BSD_SYS_EVENT_H_ + +#include + +#ifdef _DEFAULT_SOURCE + +#include + + +#define EVFILT_READ (-1) +#define EVFILT_WRITE (-2) +#define EVFILT_PROC (-5) + +#if defined(__STDC_VERSION__) && __STDC_VERSION__ >= 199901L +#define EV_SET(kevp_, a, b, c, d, e, f) do { \ + *(kevp_) = (struct kevent){ \ + .ident = (a), \ + .filter = (b), \ + .flags = (c), \ + .fflags = (d), \ + .data = (e), \ + .udata = (f), \ + .ext = {0}, \ + }; \ +} while (0) +#else /* Pre-C99 or not STDC (e.g., C++) */ +/* + * The definition of the local variable kevp could possibly conflict + * with a user-defined value passed in parameters a-f. + */ +#define EV_SET(kevp_, a, b, c, d, e, f) do { \ + struct kevent *kevp = (kevp_); \ + (kevp)->ident = (a); \ + (kevp)->filter = (b); \ + (kevp)->flags = (c); \ + (kevp)->fflags = (d); \ + (kevp)->data = (e); \ + (kevp)->udata = (f); \ + (kevp)->ext[0] = 0; \ + (kevp)->ext[1] = 0; \ + (kevp)->ext[2] = 0; \ + (kevp)->ext[3] = 0; \ +} while (0) +#endif + +struct kevent { + uintptr_t ident; /* identifier for this event */ + short filter; /* filter for event */ + unsigned short flags; /* action flags for kqueue */ + unsigned int fflags; /* filter flag value */ + int64_t data; /* filter data value */ + void *udata; /* opaque user data identifier */ + uint64_t ext[4]; /* extensions */ +}; + +/* actions */ +#define EV_ADD 0x0001 /* add event to kq (implies enable) */ +#define EV_DELETE 0x0002 /* delete event from kq */ + +/* flags */ +#define EV_ONESHOT 0x0010 /* only report one occurrence */ +#define EV_CLEAR 0x0020 /* clear event state after reporting */ + +/* returned values */ +#define EV_EOF 0x8000 /* EOF detected */ +#define EV_ERROR 0x4000 /* error, data contains errno */ + +/* data/hint flags for EVFILT_PROC */ +#define NOTE_EXIT 0x80000000 /* process exited */ + + +#ifdef __cplusplus +extern "C" { +#endif + +int kqueue(void); +int kevent(int kq, const struct kevent *changelist, int nchanges, + struct kevent *eventlist, int nevents, + const struct timespec *timeout); + +#ifdef __cplusplus +} +#endif + + +#endif /* _DEFAULT_SOURCE */ + +#endif /* _BSD_SYS_EVENT_H_ */ diff --git a/src/libs/bsd/Jamfile b/src/libs/bsd/Jamfile index c34154a01f..aa5ea6df69 100644 --- a/src/libs/bsd/Jamfile +++ b/src/libs/bsd/Jamfile @@ -1,5 +1,7 @@ SubDir HAIKU_TOP src libs bsd ; +UsePrivateHeaders system ; +UsePrivateHeaders [ FDirName system arch $(TARGET_ARCH) ] ; UseHeaders [ FDirName $(HAIKU_TOP) headers compatibility bsd ] : true ; local architectureObject ; @@ -17,6 +19,7 @@ for architectureObject in [ MultiArchSubDirSetup ] { fts.c getpass.c issetugid.c + kqueue.cpp lutimes.c progname.c pty.cpp diff --git a/src/libs/bsd/kqueue.cpp b/src/libs/bsd/kqueue.cpp new file mode 100644 index 0000000000..122b4a51bf --- /dev/null +++ b/src/libs/bsd/kqueue.cpp @@ -0,0 +1,263 @@ +/* + * Copyright 2023, Haiku, Inc. All rights reserved. + * Distributed under the terms of the MIT License. + */ +#include + +#include + +#include +#include +#include + + +extern "C" int +kqueue() +{ + int fd = _kern_event_queue_create(0); + if (fd < 0) { + __set_errno(fd); + return -1; + } + return fd; +} + + +static short +filter_from_info(const event_wait_info& info) +{ + switch (info.type) { + case B_OBJECT_TYPE_FD: + if (info.events > 0 && (info.events & B_EVENT_WRITE) != 0) + return EVFILT_WRITE; + return EVFILT_READ; + + case B_OBJECT_TYPE_THREAD: + return EVFILT_PROC; + } + + return 0; +} + + +extern "C" int +kevent(int kq, + const struct kevent *changelist, int nchanges, + struct kevent *eventlist, int nevents, + const struct timespec *tspec) +{ + BStackOrHeapArray waitInfos(max_c(nchanges, nevents)); + + event_wait_info* waitInfo = waitInfos; + int changedInfos = 0; + + for (int i = 0; i < nchanges; i++) { + waitInfo->object = changelist[i].ident; + waitInfo->events = 0; + waitInfo->user_data = changelist[i].udata; + + int32 events = 0, behavior = 0; + switch (changelist[i].filter) { + case EVFILT_READ: + waitInfo->type = B_OBJECT_TYPE_FD; + events = B_EVENT_READ; + break; + + case EVFILT_WRITE: + waitInfo->type = B_OBJECT_TYPE_FD; + events = B_EVENT_WRITE; + break; + + case EVFILT_PROC: + waitInfo->type = B_OBJECT_TYPE_THREAD; + if ((changelist[i].fflags & NOTE_EXIT) != 0) + events |= B_EVENT_INVALID; + break; + + default: + return EINVAL; + } + + if ((changelist[i].flags & EV_ONESHOT) != 0) + behavior |= B_EVENT_ONE_SHOT; + if ((changelist[i].flags & EV_CLEAR) == 0) + behavior |= B_EVENT_LEVEL_TRIGGERED; + + if (changelist[i].filter == EVFILT_READ || changelist[i].filter == EVFILT_WRITE) { + // kqueue treats the same file descriptor with both READ and WRITE filters + // as two separate listeners. Haiku, however, treats it as one. + // We rectify this here by carefully combining the two. + + // We can't support ONESHOT for descriptors due to the separation. + if ((changelist[i].flags & EV_ONESHOT) != 0) { + __set_errno(EOPNOTSUPP); + return -1; + } + + const short otherFilter = (changelist[i].filter == EVFILT_READ) + ? EVFILT_WRITE : EVFILT_READ; + const int32 otherEvents = (otherFilter == EVFILT_READ) + ? B_EVENT_READ : B_EVENT_WRITE; + + // First, check if the other filter is specified in this changelist. + int j; + for (j = 0; j < nchanges; j++) { + if (changelist[j].ident != changelist[i].ident) + continue; + if (changelist[j].filter != otherFilter) + continue; + + // We've found it. + break; + } + if (j < nchanges) { + // It is in the list. + if (j < i) { + // And it's already been taken care of. + continue; + } + + // Fold it into this one. + if ((changelist[j].flags & EV_ADD) != 0) { + waitInfo->events |= otherEvents; + } else if ((changelist[j].flags & EV_DELETE) != 0) { + waitInfo->events &= ~otherEvents; + } + } else { + // It is not in the list. See if it's already set. + event_wait_info info; + info.type = B_OBJECT_TYPE_FD; + info.object = waitInfo->object; + info.events = -1; + + status_t status = _kern_event_queue_select(kq, &info, 1); + if (status == B_OK) + waitInfo->events |= (info.events & otherEvents); + } + } + + if ((changelist[i].flags & EV_ADD) != 0) { + waitInfo->events |= events; + } else if ((changelist[i].flags & EV_DELETE) != 0) { + waitInfo->events &= ~events; + } + + if (waitInfo->events != 0) + waitInfo->events |= behavior; + + changedInfos++; + waitInfo++; + } + if (changedInfos != 0) { + status_t status = _kern_event_queue_select(kq, waitInfos, changedInfos); + if (status != B_OK) { + if (nchanges == 1 && nevents == 0) { + // Special case: return the lone error directly. + __set_errno(waitInfos[0].events); + return -1; + } + + // Report problems as error events. + int errors = 0; + for (int i = 0; i < changedInfos; i++) { + if (waitInfos[i].events > 0) + continue; + if (nevents == 0) + break; + + short filter = filter_from_info(waitInfos[i]); + int64_t data = waitInfos[i].events; + EV_SET(eventlist, waitInfos[i].object, + filter, EV_ERROR, 0, data, waitInfos[i].user_data); + eventlist++; + nevents--; + errors++; + } + if (nevents == 0 || errors == 0) { + __set_errno(status); + return -1; + } + } + } + + if (nevents != 0) { + bigtime_t timeout = 0; + uint32 waitFlags = 0; + if (tspec != NULL) { + timeout = (tspec->tv_sec * 1000000LL) + (tspec->tv_nsec / 1000LL); + waitFlags |= B_RELATIVE_TIMEOUT; + } + + ssize_t events = _kern_event_queue_wait(kq, waitInfos, + max_c(1, nevents / 2), waitFlags, timeout); + if (events > 0) { + int returnedEvents = 0; + for (ssize_t i = 0; i < events; i++) { + unsigned short flags = 0; + unsigned int fflags = 0; + int64_t data = 0; + + if (waitInfos[i].events < 0) { + flags |= EV_ERROR; + data = waitInfos[i].events; + } else if ((waitInfos[i].events & B_EVENT_DISCONNECTED) != 0) { + flags |= EV_EOF; + } else if ((waitInfos[i].events & B_EVENT_INVALID) != 0) { + switch (waitInfos[i].type) { + case B_OBJECT_TYPE_FD: + flags |= EV_EOF; + break; + + case B_OBJECT_TYPE_THREAD: { + fflags |= NOTE_EXIT; + + status_t returnValue = -1; + status_t status = wait_for_thread(waitInfos[i].object, &returnValue); + if (status == B_OK) + data = returnValue; + else + data = -1; + break; + } + } + } else if ((waitInfos[i].events & B_EVENT_ERROR) != 0) { + flags |= EV_ERROR; + data = EINVAL; + } + + short filter = filter_from_info(waitInfos[i]); + if (waitInfos[i].type == B_OBJECT_TYPE_FD && (flags & (EV_ERROR | EV_EOF)) == 0) { + // Do we have both a read and a write event? + if ((waitInfos[i].events & (B_EVENT_READ | B_EVENT_WRITE)) + == (B_EVENT_READ | B_EVENT_WRITE)) { + // We do. Report both, if we can. + if (nevents > 1) { + EV_SET(eventlist, waitInfos[i].object, + EVFILT_WRITE, flags, fflags, data, waitInfos[i].user_data); + eventlist++; + returnedEvents++; + nevents--; + } + filter = EVFILT_READ; + } + } + + EV_SET(eventlist, waitInfos[i].object, + filter, flags, fflags, data, waitInfos[i].user_data); + eventlist++; + returnedEvents++; + nevents--; + } + return returnedEvents; + } else if (events < 0) { + if (events == B_WOULD_BLOCK || events == B_TIMED_OUT) + return 0; + + __set_errno(events); + return -1; + } + return 0; + } + + return 0; +}