Groups | Search | Server Info | Keyboard shortcuts | Login | Register [http] [https] [nntp] [nntps]
Groups > comp.programming.threads > #2570 > unrolled thread
| Started by | Steven Stewart-Gallus <stevenselectronicmail@gmail.com> |
|---|---|
| First post | 2014-08-08 19:17 -0700 |
| Last post | 2014-08-10 20:57 -0700 |
| Articles | 3 — 2 participants |
Back to article view | Back to comp.programming.threads
Can I get code review for some shared memory RPC code? Steven Stewart-Gallus <stevenselectronicmail@gmail.com> - 2014-08-08 19:17 -0700
Re: Can I get code review for some shared memory RPC code? andrew@cucumber.demon.co.uk (Andrew Gabriel) - 2014-08-10 15:45 +0000
Re: Can I get code review for some shared memory RPC code? Steven Stewart-Gallus <stevenselectronicmail@gmail.com> - 2014-08-10 20:57 -0700
| From | Steven Stewart-Gallus <stevenselectronicmail@gmail.com> |
|---|---|
| Date | 2014-08-08 19:17 -0700 |
| Subject | Can I get code review for some shared memory RPC code? |
| Message-ID | <e10e88a1-5526-4257-9053-c98b49fa31ed@googlegroups.com> |
The code should be fairly portable (at least with futex use disabled).
Of course, I'm using POSIX threads as C11 threads aren't widely
implemented. I have absolutely zero experience doing any sort of
complicated multithreaded programming with atomics (I prefer to stick
to isolated heaps and very simple message queues) so the code is
probably broken in some way. One thing that doesn't help is that I
can only test on x86 right now (which has a very strong memory model).
Annoyingly, there's no portable pause intrinsic but it should be
really easy to s/_mm_pause/my_pause/g for a different architecture.
#ifdef USE_FUTEX
#define _GNU_SOURCE
#endif
#include <assert.h>
#include <errno.h>
#include <inttypes.h>
#include <pthread.h>
#include <stdatomic.h>
#include <stdbool.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/syscall.h>
#ifdef USE_FUTEX
#include <unistd.h>
#include <linux/futex.h>
#endif
#include <immintrin.h>
#ifndef __STDC_LIB_EXT1__
#define errno_t int
#endif
struct ssg_rpc_callsite {
atomic_int request_pending;
atomic_int reply_pending;
};
errno_t ssg_rpc_recv_request(struct ssg_rpc_callsite * callsite);
errno_t ssg_rpc_send_reply(struct ssg_rpc_callsite * callsite);
errno_t ssg_rpc_send_request(struct ssg_rpc_callsite * callsite);
errno_t ssg_rpc_recv_reply(struct ssg_rpc_callsite * callsite);
struct request {
unsigned data;
};
struct reply {
unsigned data;
};
union request_or_reply {
struct request request;
struct reply reply;
};
struct shm {
struct ssg_rpc_callsite callsite;
union request_or_reply data;
};
static struct {
struct shm shm;
char _padding0[128U - sizeof (struct shm) % 128U];
} conn_mem;
static struct {
struct shm shm;
char _padding0[128U - sizeof (struct shm) % 128U];
} server_mem;
static struct {
struct shm shm;
char _padding0[128U - sizeof (struct shm) % 128U];
} client_mem;
static void * server_task(void *arg);
static void * server_friend_task(void *arg);
static void * client_task(void *arg);
static void * client_friend_task(void *arg);
/* In a real situation the code would probably be used in a
* multiprocess situation where the processes don't trust each other.
* Currently, I am using threads and not processes so I can use cool
* tools like thread sanitizer to debug the code. However, to
* accurately simulate the performance of the multiprocess situation
* where the tasks don't trust each other I am using more threads than
* you'd need in a situation when the tasks trust each other.
* Multiple processes would be needed when the communicators don't
* trust each other in order to deal with with attackers maliciously
* resizing the shared memory file used to communicate.
*/
int main(void)
{
errno_t errnum;
pthread_t server_thread;
{
pthread_t xx;
if ((errnum = pthread_create(&xx, NULL, server_task, NULL)) != 0) {
errno = errnum;
perror("pthread_create");
return EXIT_FAILURE;
}
server_thread = xx;
}
pthread_t client_thread;
{
pthread_t xx;
if ((errnum = pthread_create(&xx, NULL, client_task, NULL)) != 0) {
errno = errnum;
perror("pthread_create");
return EXIT_FAILURE;
}
client_thread = xx;
}
pthread_join(client_thread, NULL);
return EXIT_SUCCESS;
}
static void * server_task(void *arg)
{
errno_t errnum;
pthread_t friend_thread;
{
pthread_t xx;
if ((errnum = pthread_create(&xx, NULL, server_friend_task, NULL)) != 0) {
errno = errnum;
perror("pthread_create");
exit(EXIT_FAILURE);
}
friend_thread = xx;
}
/* Implement the server logic with trusted data only */
for (;;) {
if ((errnum = ssg_rpc_recv_request(&server_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_recv_request");
exit(EXIT_FAILURE);
}
server_mem.shm.data.reply.data = server_mem.shm.data.request.data + 2U;
if ((errnum = ssg_rpc_send_reply(&server_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_send_reply");
exit(EXIT_FAILURE);
}
}
return NULL;
}
static void * server_friend_task(void * arg)
{
errno_t errnum;
/* Use a broker process as a sacrifice to malicious clients */
for (;;) {
if ((errnum = ssg_rpc_recv_request(&conn_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_recv_request");
exit(EXIT_FAILURE);
}
memcpy(&server_mem.shm.data.request, &conn_mem.shm.data.request,
sizeof conn_mem.shm.data.request);
if ((errnum = ssg_rpc_send_request(&server_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_send_request");
exit(EXIT_FAILURE);
}
if ((errnum = ssg_rpc_recv_reply(&server_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_recv_reply");
exit(EXIT_FAILURE);
}
memcpy(&conn_mem.shm.data.reply, &server_mem.shm.data.reply,
sizeof conn_mem.shm.data.reply);
if ((errnum = ssg_rpc_send_reply(&conn_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_send_reply");
exit(EXIT_FAILURE);
}
}
}
static void * client_task(void* arg)
{
errno_t errnum;
pthread_t friend_thread;
{
pthread_t xx;
if ((errnum = pthread_create(&xx, NULL, client_friend_task, NULL)) != 0) {
errno = errnum;
perror("pthread_create");
exit(EXIT_FAILURE);
}
friend_thread = xx;
}
/* Implement the client logic with trusted data only */
for (unsigned ii = 0U; ii < 900000U; ++ii) {
unsigned request = ii;
client_mem.shm.data.request.data = request;
if ((errnum = ssg_rpc_send_request(&client_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_send_request");
exit(EXIT_FAILURE);
}
/* Could do other stuff concurrently */
if ((errnum = ssg_rpc_recv_reply(&client_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_recv_reply");
exit(EXIT_FAILURE);
}
unsigned reply = client_mem.shm.data.reply.data;
if (reply != request + 2U) {
printf("request + 2: %u, reply: %u\n", request + 2U, reply);
exit(EXIT_FAILURE);
}
}
client_mem.shm.data.request.data = 3U;
printf("requesting: %i\n", client_mem.shm.data.request.data);
if ((errnum = ssg_rpc_send_request(&client_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_send_request");
exit(EXIT_FAILURE);
}
if ((errnum = ssg_rpc_recv_reply(&client_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_recv_reply");
exit(EXIT_FAILURE);
}
printf("was replied with: %i\n", client_mem.shm.data.reply.data);
return NULL;
}
static void *client_friend_task(void *arg)
{
errno_t errnum;
/* Use a broker process as a sacrifice to malicious servers */
for (;;) {
if ((errnum = ssg_rpc_recv_request(&client_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_recv_request");
exit(EXIT_FAILURE);
}
memcpy(&conn_mem.shm.data.request, &client_mem.shm.data.request,
sizeof conn_mem.shm.data.request);
if ((errnum = ssg_rpc_send_request(&conn_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_send_request");
exit(EXIT_FAILURE);
}
if ((errnum = ssg_rpc_recv_reply(&conn_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_recv_reply");
exit(EXIT_FAILURE);
}
memcpy(&client_mem.shm.data.reply, &conn_mem.shm.data.reply,
sizeof conn_mem.shm.data.reply);
if ((errnum = ssg_rpc_send_reply(&client_mem.shm.callsite)) != 0) {
errno = errnum;
perror("ssg_rpc_send_reply");
exit(EXIT_FAILURE);
}
}
}
#define RECV_SPIN_COUNT 800U
#define SEND_SPIN_COUNT 200U
static errno_t wait_until_different(atomic_int const *uaddr, int val);
static errno_t hint_wakeup(atomic_int const *uaddr);
errno_t ssg_rpc_recv_request(struct ssg_rpc_callsite * callsite)
{
errno_t errnum;
for (;;) {
for (size_t ii = 0U; ii < RECV_SPIN_COUNT; ++ii) {
bool request_pending = atomic_load_explicit(&callsite->request_pending,
memory_order_acquire);
if (request_pending) {
goto got_reply;
}
_mm_pause();
}
errnum = wait_until_different(&callsite->request_pending, false);
switch (errnum) {
case EAGAIN:
goto got_reply;
case 0:
case EINTR:
continue;
default:
return errnum;
}
}
got_reply:
atomic_store_explicit(&callsite->request_pending, false,
memory_order_release);
atomic_thread_fence(memory_order_acquire);
return 0;
}
errno_t ssg_rpc_send_reply(struct ssg_rpc_callsite * callsite)
{
atomic_thread_fence(memory_order_release);
atomic_store_explicit(&callsite->reply_pending, true, memory_order_release);
for (size_t ii = 0U; ii < SEND_SPIN_COUNT; ++ii) {
_mm_pause();
bool reply_pending = atomic_load_explicit(&callsite->reply_pending,
memory_order_acquire);
if (!reply_pending) {
return 0;
}
}
return hint_wakeup(&callsite->reply_pending);
}
errno_t ssg_rpc_send_request(struct ssg_rpc_callsite * callsite)
{
atomic_thread_fence(memory_order_release);
atomic_store_explicit(&callsite->request_pending, true,
memory_order_release);
for (size_t ii = 0U; ii < SEND_SPIN_COUNT; ++ii) {
_mm_pause();
bool request_pending = atomic_load_explicit(&callsite->request_pending,
memory_order_acquire);
if (!request_pending) {
return 0;
}
}
return hint_wakeup(&callsite->request_pending);
}
errno_t ssg_rpc_recv_reply(struct ssg_rpc_callsite * callsite)
{
errno_t errnum;
for (;;) {
for (size_t ii = 0U; ii < RECV_SPIN_COUNT; ++ii) {
bool reply_pending = atomic_load_explicit(&callsite->reply_pending,
memory_order_acquire);
if (reply_pending) {
goto got_reply;
}
_mm_pause();
}
errnum = wait_until_different(&callsite->reply_pending, false);
switch (errnum) {
case EAGAIN:
goto got_reply;
case 0:
case EINTR:
continue;
default:
return errnum;
}
}
got_reply:
atomic_store_explicit(&callsite->reply_pending, false,
memory_order_release);
atomic_thread_fence(memory_order_acquire);
return 0;
}
#ifdef USE_FUTEX
static errno_t futex_wait(atomic_int const *uaddr, int val, struct timespec const *timeout);
static errno_t futex_wake(unsigned * restrict wokeupp, atomic_int const *uaddr, int val);
static errno_t wait_until_different(atomic_int const *uaddr, int val)
{
return futex_wait(uaddr, val, NULL);
}
static errno_t hint_wakeup(atomic_int const *uaddr)
{
return futex_wake(NULL, uaddr, 1);
}
_Static_assert(ATOMIC_INT_LOCK_FREE == 2,
"lock using atomic_uints are incompatible with futexes");
static errno_t futex_wait(atomic_int const *uaddr, int val, struct timespec const *timeout)
{
int xx = syscall(__NR_futex, (intptr_t)uaddr, (intptr_t)FUTEX_WAIT,
(intptr_t)val,
(intptr_t)timeout);
if (xx < 0) {
return errno;
}
return 0;
}
static errno_t futex_wake(unsigned * restrict wokeupp, atomic_int const *uaddr,
int val)
{
int xx = syscall(__NR_futex, (intptr_t)uaddr, (intptr_t)FUTEX_WAKE,
(intptr_t)val);
if (xx < 0) {
return errno;
}
if (wokeupp != NULL) {
*wokeupp = xx;
}
return 0;
}
#else
static errno_t wait_until_different(atomic_int const *uaddr, int val)
{
if (atomic_load_explicit(uaddr, memory_order_acquire) != val) {
return EAGAIN;
}
for (;;) {
_mm_pause();
if (atomic_load_explicit(uaddr, memory_order_acquire) != val) {
return 0;
}
}
}
static errno_t hint_wakeup(atomic_int const *uaddr)
{
return 0;
}
#endif
[toc] | [next] | [standalone]
| From | andrew@cucumber.demon.co.uk (Andrew Gabriel) |
|---|---|
| Date | 2014-08-10 15:45 +0000 |
| Message-ID | <ls843m$945$1@dont-email.me> |
| In reply to | #2570 |
In article <e10e88a1-5526-4257-9053-c98b49fa31ed@googlegroups.com>, Steven Stewart-Gallus <stevenselectronicmail@gmail.com> writes: > The code should be fairly portable (at least with futex use disabled). > Of course, I'm using POSIX threads as C11 threads aren't widely > implemented. I have absolutely zero experience doing any sort of > complicated multithreaded programming with atomics (I prefer to stick > to isolated heaps and very simple message queues) so the code is > probably broken in some way. One thing that doesn't help is that I > can only test on x86 right now (which has a very strong memory model). > Annoyingly, there's no portable pause intrinsic but it should be > really easy to s/_mm_pause/my_pause/g for a different architecture. Your wait_until_different and hint_wakeup should be replaced by a condition variable (cond_wait, and cond_signal or cond_broadcast). Then the need for your pause function and spin loop go away. I don't know how portable the atomic_* functions are that you are using. I would use a mutex protected variable, which should work everywhere, if portability is important. I did not look at the program logic in detail, as you didn't describe what you intended it to do. -- Andrew Gabriel [email address is not usable -- followup in the newsgroup]
[toc] | [prev] | [next] | [standalone]
| From | Steven Stewart-Gallus <stevenselectronicmail@gmail.com> |
|---|---|
| Date | 2014-08-10 20:57 -0700 |
| Message-ID | <45f15ddc-9fdc-4e0e-a2b6-2e72e2bd52ed@googlegroups.com> |
| In reply to | #2571 |
> Your wait_until_different and hint_wakeup should be replaced by a > condition variable (cond_wait, and cond_signal or cond_broadcast). > Then the need for your pause function and spin loop go away. Condition variables are implemented in terms of the futex system call (on Linux only obviously). I am reimplementing my own concurrency primitives in terms of the lowest level OS primitives for performance reasons and for learning experience. If I used the condition variables provided by GLibc it'd be highly likely that I'd still be using spin loops, they'd just be hidden inside of GLibc. As well, when built with futexes enabled the code only spins for a small amount of time and waits in a kernel wait queue otherwise. I'm honestly kind of offended because you seemed to have posted a very silly comment without taking the time to learn a bit about the type of code you were reviewing first. > I don't know how portable the atomic_* functions are that you are > using. I would use a mutex protected variable, which should work > everywhere, if portability is important. stdatomic.h is part of the C11 standard's new functionality for multithreading. I'm honestly surprised that a poster on comp.programming.threads doesn't know about stdatomic.h or can't bother to look it up. The futex system call is obviously unportable but it is hidden behind a define. _mm_pause is an Intel intrinsic and is obviously nonportable but it should be easy to find a platform equivalent for any one that I need to port to. It's very amusing that you go on to recommend to use a mutex protected variable as only the atomics part of the C11 standard is very widely implemented, and POSIX mutexes and condition variables are not portable to Windows. Moreover, POSIX mutexes and condition variables are unlikely to be portable to embedded platforms or other specialized hardware. atomic_* functions are probably MORE portable than using any libraries implementations of mutex. Besides, for performance I need my own custom implementation in terms of the lowest level primitives. > I did not look at the program logic in detail, as you didn't describe > what you intended it to do. What's left to say? I'm implementing Shared Memory Remote Procedure Calls between two tasks.
[toc] | [prev] | [standalone]
Back to top | Article view | comp.programming.threads
csiph-web