mirror of
https://github.com/wassname/ray.git
synced 2026-08-09 12:20:09 +08:00
Run clang-format and check in Travis CI (#14)
* Run clang-format and add pre-commit hook for it. * Modify .travis.yml to check * Try to fix problems with .travis.yml * Try to fix .travis.yml yet again * Update .clang-format to Philipp's preferences * Don't allow lint to fail in Travis * Remove git-hooks directory * Improve clang-format failure output * Fix clang-format error * Report which commit clang-format is comparing against, and add whitespace error * Handle non-PR Travis in clang-format, and add another error * Check $TRAVIS_PULL_REQUEST correctly and add another error * Fix syntax error in check-git-clang-format-output.sh * Add whitespace error * Remove extra whitespace, add clang-format to README
This commit is contained in:
committed by
Robert Nishihara
parent
a62c0f8fac
commit
04737f3f56
+29
-17
@@ -3,8 +3,8 @@
|
||||
#include <assert.h>
|
||||
#include <unistd.h>
|
||||
|
||||
UT_icd item_icd = { sizeof(event_loop_item), NULL, NULL, NULL };
|
||||
UT_icd poll_icd = { sizeof(struct pollfd), NULL, NULL, NULL };
|
||||
UT_icd item_icd = {sizeof(event_loop_item), NULL, NULL, NULL};
|
||||
UT_icd poll_icd = {sizeof(struct pollfd), NULL, NULL, NULL};
|
||||
|
||||
/* Initializes the event loop.
|
||||
* This function needs to be called before any other event loop function. */
|
||||
@@ -18,15 +18,19 @@ void event_loop_init(event_loop *loop) {
|
||||
* which can be queried using event_loop_type and event_loop_id. The parameter
|
||||
* events is the same as in http://linux.die.net/man/2/poll.
|
||||
* Returns the index of the item in the event loop. */
|
||||
int64_t event_loop_attach(event_loop *loop, int type, data_connection* connection, int fd, int events) {
|
||||
int64_t event_loop_attach(event_loop *loop,
|
||||
int type,
|
||||
data_connection *connection,
|
||||
int fd,
|
||||
int events) {
|
||||
assert(utarray_len(loop->items) == utarray_len(loop->waiting));
|
||||
int64_t index = utarray_len(loop->items);
|
||||
event_loop_item item = { .type = type };
|
||||
event_loop_item item = {.type = type};
|
||||
if (connection) {
|
||||
item.connection = *connection;
|
||||
}
|
||||
utarray_push_back(loop->items, &item );
|
||||
struct pollfd waiting = { .fd = fd, .events = events };
|
||||
utarray_push_back(loop->items, &item);
|
||||
struct pollfd waiting = {.fd = fd, .events = events};
|
||||
utarray_push_back(loop->waiting, &waiting);
|
||||
return index;
|
||||
}
|
||||
@@ -35,16 +39,18 @@ int64_t event_loop_attach(event_loop *loop, int type, data_connection* connectio
|
||||
* This invalidates all other indices into the event loop items, but leaves
|
||||
* the ids of the event loop items valid. */
|
||||
void event_loop_detach(event_loop *loop, int64_t index, int shall_close) {
|
||||
struct pollfd *waiting_item = (struct pollfd*) utarray_eltptr(loop->waiting, index);
|
||||
struct pollfd *waiting_back = (struct pollfd*) utarray_back(loop->waiting);
|
||||
struct pollfd *waiting_item =
|
||||
(struct pollfd *) utarray_eltptr(loop->waiting, index);
|
||||
struct pollfd *waiting_back = (struct pollfd *) utarray_back(loop->waiting);
|
||||
if (shall_close) {
|
||||
close(waiting_item->fd);
|
||||
}
|
||||
*waiting_item = *waiting_back;
|
||||
utarray_pop_back(loop->waiting);
|
||||
|
||||
event_loop_item *items_item = (event_loop_item*) utarray_eltptr(loop->items, index);
|
||||
event_loop_item *items_back = (event_loop_item*) utarray_back(loop->items);
|
||||
event_loop_item *items_item =
|
||||
(event_loop_item *) utarray_eltptr(loop->items, index);
|
||||
event_loop_item *items_back = (event_loop_item *) utarray_back(loop->items);
|
||||
*items_item = *items_back;
|
||||
utarray_pop_back(loop->items);
|
||||
}
|
||||
@@ -52,7 +58,8 @@ void event_loop_detach(event_loop *loop, int64_t index, int shall_close) {
|
||||
/* Poll the file descriptors associated to this event loop.
|
||||
* See http://linux.die.net/man/2/poll */
|
||||
int event_loop_poll(event_loop *loop) {
|
||||
return poll((struct pollfd*) utarray_front(loop->waiting), utarray_len(loop->waiting), -1);
|
||||
return poll((struct pollfd *) utarray_front(loop->waiting),
|
||||
utarray_len(loop->waiting), -1);
|
||||
}
|
||||
|
||||
/* Get the total number of file descriptors participating in the event loop. */
|
||||
@@ -60,20 +67,25 @@ int64_t event_loop_size(event_loop *loop) {
|
||||
return utarray_len(loop->waiting);
|
||||
}
|
||||
|
||||
/* Get the pollfd structure associated to a file descriptor participating in the event loop. */
|
||||
/* Get the pollfd structure associated to a file descriptor participating in the
|
||||
* event loop. */
|
||||
struct pollfd *event_loop_get(event_loop *loop, int64_t index) {
|
||||
return (struct pollfd*) utarray_eltptr(loop->waiting, index);
|
||||
return (struct pollfd *) utarray_eltptr(loop->waiting, index);
|
||||
}
|
||||
|
||||
/* Set the data connection information for participant in the event loop. */
|
||||
void event_loop_set_connection(event_loop *loop, int64_t index, const data_connection* conn) {
|
||||
event_loop_item *item = (event_loop_item*) utarray_eltptr(loop->items, index);
|
||||
void event_loop_set_connection(event_loop *loop,
|
||||
int64_t index,
|
||||
const data_connection *conn) {
|
||||
event_loop_item *item =
|
||||
(event_loop_item *) utarray_eltptr(loop->items, index);
|
||||
item->connection = *conn;
|
||||
}
|
||||
|
||||
/* Get the data connection information for participant in the event loop. */
|
||||
data_connection* event_loop_get_connection(event_loop *loop, int64_t index) {
|
||||
event_loop_item *item = (event_loop_item*) utarray_eltptr(loop->items, index);
|
||||
data_connection *event_loop_get_connection(event_loop *loop, int64_t index) {
|
||||
event_loop_item *item =
|
||||
(event_loop_item *) utarray_eltptr(loop->items, index);
|
||||
return &item->connection;
|
||||
}
|
||||
|
||||
|
||||
+11
-5
@@ -12,25 +12,31 @@ typedef struct {
|
||||
int type;
|
||||
/* If type is data transfer, this contains information about the status
|
||||
* of the transfer. */
|
||||
data_connection connection;
|
||||
data_connection connection;
|
||||
} event_loop_item;
|
||||
|
||||
typedef struct {
|
||||
/* Array of event_loop_items that hold information for connections. */
|
||||
UT_array *items;
|
||||
UT_array *items;
|
||||
/* Array of file descriptors that are waiting, corresponding to items. */
|
||||
UT_array *waiting;
|
||||
UT_array *waiting;
|
||||
} event_loop;
|
||||
|
||||
/* Event loop functions. */
|
||||
void event_loop_init(event_loop *loop);
|
||||
void event_loop_free(event_loop *loop);
|
||||
int64_t event_loop_attach(event_loop *loop, int type, data_connection* connection, int fd, int events);
|
||||
int64_t event_loop_attach(event_loop *loop,
|
||||
int type,
|
||||
data_connection *connection,
|
||||
int fd,
|
||||
int events);
|
||||
void event_loop_detach(event_loop *loop, int64_t index, int shall_close);
|
||||
int event_loop_poll(event_loop *loop);
|
||||
int64_t event_loop_size(event_loop *loop);
|
||||
struct pollfd *event_loop_get(event_loop *loop, int64_t index);
|
||||
void event_loop_set_connection(event_loop *loop, int64_t index, const data_connection* conn);
|
||||
void event_loop_set_connection(event_loop *loop,
|
||||
int64_t index,
|
||||
const data_connection *conn);
|
||||
data_connection *event_loop_get_connection(event_loop *loop, int64_t index);
|
||||
|
||||
#endif
|
||||
|
||||
+4
-4
@@ -1,7 +1,7 @@
|
||||
/* A simple example on how to use the plasma store
|
||||
*
|
||||
*
|
||||
* Can be called in the following way:
|
||||
*
|
||||
*
|
||||
* cd build
|
||||
* ./plasma_store -s /tmp/plasma_socket
|
||||
* ./example -s /tmp/plasma_socket -g
|
||||
@@ -19,8 +19,8 @@ int main(int argc, char *argv[]) {
|
||||
int64_t size;
|
||||
void *data;
|
||||
int c;
|
||||
plasma_id id = {{255, 255, 255, 255, 255, 255, 255, 255, 255, 255, 255, 255,
|
||||
255, 255, 255, 255, 255, 255, 255, 255}};
|
||||
plasma_id id = {{255, 255, 255, 255, 255, 255, 255, 255, 255, 255,
|
||||
255, 255, 255, 255, 255, 255, 255, 255, 255, 255}};
|
||||
while ((c = getopt(argc, argv, "s:cfg")) != -1) {
|
||||
switch (c) {
|
||||
case 's':
|
||||
|
||||
+13
-8
@@ -1,7 +1,9 @@
|
||||
#include "fling.h"
|
||||
|
||||
void init_msg(struct msghdr *msg, struct iovec *iov,
|
||||
char *buf, size_t buf_len) {
|
||||
void init_msg(struct msghdr *msg,
|
||||
struct iovec *iov,
|
||||
char *buf,
|
||||
size_t buf_len) {
|
||||
iov->iov_base = buf;
|
||||
iov->iov_len = 1;
|
||||
|
||||
@@ -13,7 +15,7 @@ void init_msg(struct msghdr *msg, struct iovec *iov,
|
||||
msg->msg_namelen = 0;
|
||||
}
|
||||
|
||||
int send_fd(int conn, int fd, const char* payload, int size) {
|
||||
int send_fd(int conn, int fd, const char *payload, int size) {
|
||||
struct msghdr msg;
|
||||
struct iovec iov;
|
||||
char buf[CMSG_SPACE(sizeof(int))];
|
||||
@@ -24,13 +26,13 @@ int send_fd(int conn, int fd, const char* payload, int size) {
|
||||
header->cmsg_level = SOL_SOCKET;
|
||||
header->cmsg_type = SCM_RIGHTS;
|
||||
header->cmsg_len = CMSG_LEN(sizeof(int));
|
||||
*(int *)CMSG_DATA(header) = fd;
|
||||
*(int *) CMSG_DATA(header) = fd;
|
||||
|
||||
/* send file descriptor and payload */
|
||||
return sendmsg(conn, &msg, 0) != -1 && send(conn, payload, size, 0) == -1;
|
||||
}
|
||||
|
||||
int recv_fd(int conn, char* payload, int size) {
|
||||
int recv_fd(int conn, char *payload, int size) {
|
||||
struct msghdr msg;
|
||||
struct iovec iov;
|
||||
char buf[CMSG_SPACE(sizeof(int))];
|
||||
@@ -41,11 +43,14 @@ int recv_fd(int conn, char* payload, int size) {
|
||||
|
||||
int found_fd = -1;
|
||||
int oh_noes = 0;
|
||||
for (struct cmsghdr *header = CMSG_FIRSTHDR(&msg); header != NULL; header = CMSG_NXTHDR(&msg, header))
|
||||
for (struct cmsghdr *header = CMSG_FIRSTHDR(&msg); header != NULL;
|
||||
header = CMSG_NXTHDR(&msg, header))
|
||||
if (header->cmsg_level == SOL_SOCKET && header->cmsg_type == SCM_RIGHTS) {
|
||||
int count = (header->cmsg_len - (CMSG_DATA(header) - (unsigned char *)header)) / sizeof(int);
|
||||
int count =
|
||||
(header->cmsg_len - (CMSG_DATA(header) - (unsigned char *) header)) /
|
||||
sizeof(int);
|
||||
for (int i = 0; i < count; ++i) {
|
||||
int fd = ((int *)CMSG_DATA(header))[i];
|
||||
int fd = ((int *) CMSG_DATA(header))[i];
|
||||
if (found_fd == -1) {
|
||||
found_fd = fd;
|
||||
} else {
|
||||
|
||||
+8
-7
@@ -15,20 +15,21 @@
|
||||
#include <sys/socket.h>
|
||||
#include <sys/un.h>
|
||||
|
||||
/* This is neccessary for Mac OS X, see http://www.apuebook.com/faqs2e.html (10). */
|
||||
/* This is neccessary for Mac OS X, see http://www.apuebook.com/faqs2e.html
|
||||
* (10). */
|
||||
#if !defined(CMSG_SPACE) && !defined(CMSG_LEN)
|
||||
#define CMSG_SPACE(len) (__DARWIN_ALIGN32(sizeof(struct cmsghdr)) + __DARWIN_ALIGN32(len))
|
||||
#define CMSG_LEN(len) (__DARWIN_ALIGN32(sizeof(struct cmsghdr)) + (len))
|
||||
#define CMSG_SPACE(len) \
|
||||
(__DARWIN_ALIGN32(sizeof(struct cmsghdr)) + __DARWIN_ALIGN32(len))
|
||||
#define CMSG_LEN(len) (__DARWIN_ALIGN32(sizeof(struct cmsghdr)) + (len))
|
||||
#endif
|
||||
|
||||
void init_msg(struct msghdr *msg, struct iovec *iov,
|
||||
char *buf, size_t buf_len);
|
||||
void init_msg(struct msghdr *msg, struct iovec *iov, char *buf, size_t buf_len);
|
||||
|
||||
/* Send a file descriptor "fd" and a payload "payload" of size "size"
|
||||
* over the socket "conn". Return 0 on success. */
|
||||
int send_fd(int conn, int fd, const char* payload, int size);
|
||||
int send_fd(int conn, int fd, const char *payload, int size);
|
||||
|
||||
/* Receive a file descriptor and a payload of size up to "size" from a
|
||||
* socket "conn". The payload will be written to "payload" and the file
|
||||
* descriptor will be returned. Returns -1 on failure. */
|
||||
int recv_fd(int conn, char* payload, int size);
|
||||
int recv_fd(int conn, char *payload, int size);
|
||||
|
||||
+9
-11
@@ -7,15 +7,15 @@
|
||||
#include <string.h>
|
||||
|
||||
#ifdef NDEBUG
|
||||
#define LOG_DEBUG(M, ...)
|
||||
#define LOG_DEBUG(M, ...)
|
||||
#else
|
||||
#define LOG_DEBUG(M, ...) \
|
||||
fprintf(stderr, "[DEBUG] (%s:%d) " M "\n", __FILE__, __LINE__, ##__VA_ARGS__)
|
||||
#define LOG_DEBUG(M, ...) \
|
||||
fprintf(stderr, "[DEBUG] (%s:%d) " M "\n", __FILE__, __LINE__, ##__VA_ARGS__)
|
||||
#endif
|
||||
|
||||
#define LOG_ERR(M, ...) \
|
||||
fprintf(stderr, "[ERROR] (%s:%d: errno: %s) " M "\n", \
|
||||
__FILE__, __LINE__, errno == 0 ? "None" : strerror(errno), ##__VA_ARGS__)
|
||||
#define LOG_ERR(M, ...) \
|
||||
fprintf(stderr, "[ERROR] (%s:%d: errno: %s) " M "\n", __FILE__, __LINE__, \
|
||||
errno == 0 ? "None" : strerror(errno), ##__VA_ARGS__)
|
||||
|
||||
#define LOG_INFO(M, ...) \
|
||||
fprintf(stderr, "[INFO] (%s:%d) " M "\n", __FILE__, __LINE__, ##__VA_ARGS__)
|
||||
@@ -27,9 +27,7 @@ typedef struct {
|
||||
} plasma_object_info;
|
||||
|
||||
/* Represents an object id hash, can hold a full SHA1 hash */
|
||||
typedef struct {
|
||||
unsigned char id[20];
|
||||
} plasma_id;
|
||||
typedef struct { unsigned char id[20]; } plasma_id;
|
||||
|
||||
enum plasma_request_type {
|
||||
/* Create a new object. */
|
||||
@@ -72,10 +70,10 @@ typedef struct {
|
||||
} plasma_buffer;
|
||||
|
||||
/* Connect to the local plasma store UNIX domain socket */
|
||||
int plasma_store_connect(const char* socket_name);
|
||||
int plasma_store_connect(const char *socket_name);
|
||||
|
||||
/* Connect to a possibly remote plasma manager */
|
||||
int plasma_manager_connect(const char* addr, int port);
|
||||
int plasma_manager_connect(const char *addr, int port);
|
||||
|
||||
void plasma_create(int store, plasma_id object_id, int64_t size, void **data);
|
||||
void plasma_get(int store, plasma_id object_id, int64_t *size, void **data);
|
||||
|
||||
+24
-17
@@ -25,10 +25,11 @@ void plasma_send(int fd, plasma_request *req) {
|
||||
|
||||
void plasma_create(int conn, plasma_id object_id, int64_t size, void **data) {
|
||||
LOG_INFO("called plasma_create on conn %d with size %" PRId64, conn, size);
|
||||
plasma_request req = { .type = PLASMA_CREATE, .object_id = object_id, .size = size };
|
||||
plasma_request req = {
|
||||
.type = PLASMA_CREATE, .object_id = object_id, .size = size};
|
||||
plasma_send(conn, &req);
|
||||
plasma_reply reply;
|
||||
int fd = recv_fd(conn, (char*)&reply, sizeof(plasma_reply));
|
||||
int fd = recv_fd(conn, (char *) &reply, sizeof(plasma_reply));
|
||||
assert(reply.type == PLASMA_OBJECT);
|
||||
assert(reply.size == size);
|
||||
*data = mmap(NULL, reply.size, PROT_READ | PROT_WRITE, MAP_SHARED, fd, 0);
|
||||
@@ -39,19 +40,19 @@ void plasma_create(int conn, plasma_id object_id, int64_t size, void **data) {
|
||||
}
|
||||
|
||||
void plasma_get(int conn, plasma_id object_id, int64_t *size, void **data) {
|
||||
plasma_request req = { .type = PLASMA_GET, .object_id = object_id };
|
||||
plasma_request req = {.type = PLASMA_GET, .object_id = object_id};
|
||||
plasma_send(conn, &req);
|
||||
plasma_reply reply;
|
||||
/* The following loop is run at most twice. */
|
||||
int fd = recv_fd(conn, (char*)&reply, sizeof(plasma_reply));
|
||||
int fd = recv_fd(conn, (char *) &reply, sizeof(plasma_reply));
|
||||
if (reply.type == PLASMA_FUTURE) {
|
||||
int new_fd = recv_fd(fd, (char*)&reply, sizeof(plasma_reply));
|
||||
int new_fd = recv_fd(fd, (char *) &reply, sizeof(plasma_reply));
|
||||
close(fd);
|
||||
fd = new_fd;
|
||||
}
|
||||
assert(reply.type == PLASMA_OBJECT);
|
||||
*data = mmap(NULL, reply.size, PROT_READ, MAP_SHARED, fd, 0);
|
||||
if (*data == MAP_FAILED) {
|
||||
if (*data == MAP_FAILED) {
|
||||
LOG_ERR("mmap failed");
|
||||
exit(-1);
|
||||
}
|
||||
@@ -59,11 +60,11 @@ void plasma_get(int conn, plasma_id object_id, int64_t *size, void **data) {
|
||||
}
|
||||
|
||||
void plasma_seal(int fd, plasma_id object_id) {
|
||||
plasma_request req = { .type = PLASMA_SEAL, .object_id = object_id };
|
||||
plasma_request req = {.type = PLASMA_SEAL, .object_id = object_id};
|
||||
plasma_send(fd, &req);
|
||||
}
|
||||
|
||||
int plasma_store_connect(const char* socket_name) {
|
||||
int plasma_store_connect(const char *socket_name) {
|
||||
assert(socket_name);
|
||||
struct sockaddr_un addr;
|
||||
int fd;
|
||||
@@ -73,11 +74,12 @@ int plasma_store_connect(const char* socket_name) {
|
||||
}
|
||||
memset(&addr, 0, sizeof(addr));
|
||||
addr.sun_family = AF_UNIX;
|
||||
strncpy(addr.sun_path, socket_name, sizeof(addr.sun_path)-1);
|
||||
/* Try to connect to the Plasma store. If unsuccessful, retry several times. */
|
||||
strncpy(addr.sun_path, socket_name, sizeof(addr.sun_path) - 1);
|
||||
/* Try to connect to the Plasma store. If unsuccessful, retry several times.
|
||||
*/
|
||||
int connected_successfully = 0;
|
||||
for (int num_attempts = 0; num_attempts < 50; ++num_attempts) {
|
||||
if (connect(fd, (struct sockaddr*)&addr, sizeof(addr)) == 0) {
|
||||
if (connect(fd, (struct sockaddr *) &addr, sizeof(addr)) == 0) {
|
||||
connected_successfully = 1;
|
||||
break;
|
||||
}
|
||||
@@ -94,7 +96,7 @@ int plasma_store_connect(const char* socket_name) {
|
||||
|
||||
#define h_addr h_addr_list[0]
|
||||
|
||||
int plasma_manager_connect(const char* ip_addr, int port) {
|
||||
int plasma_manager_connect(const char *ip_addr, int port) {
|
||||
int fd = socket(PF_INET, SOCK_STREAM, 0);
|
||||
if (fd < 0) {
|
||||
LOG_ERR("could not create socket");
|
||||
@@ -112,17 +114,22 @@ int plasma_manager_connect(const char* ip_addr, int port) {
|
||||
bcopy(manager->h_addr, &addr.sin_addr.s_addr, manager->h_length);
|
||||
addr.sin_port = htons(port);
|
||||
|
||||
int r = connect(fd, (struct sockaddr*) &addr, sizeof(addr));
|
||||
int r = connect(fd, (struct sockaddr *) &addr, sizeof(addr));
|
||||
if (r < 0) {
|
||||
LOG_ERR("could not establish connection to manager with id %s:%d", &ip_addr[0], port);
|
||||
LOG_ERR("could not establish connection to manager with id %s:%d",
|
||||
&ip_addr[0], port);
|
||||
exit(-1);
|
||||
}
|
||||
return fd;
|
||||
}
|
||||
|
||||
void plasma_transfer(int manager, const char* addr, int port, plasma_id object_id) {
|
||||
plasma_request req = {.type = PLASMA_TRANSFER, .object_id = object_id, .port = port};
|
||||
char* end = NULL;
|
||||
void plasma_transfer(int manager,
|
||||
const char *addr,
|
||||
int port,
|
||||
plasma_id object_id) {
|
||||
plasma_request req = {
|
||||
.type = PLASMA_TRANSFER, .object_id = object_id, .port = port};
|
||||
char *end = NULL;
|
||||
for (int i = 0; i < 4; ++i) {
|
||||
req.addr[i] = strtol(end ? end : addr, &end, 10);
|
||||
/* skip the '.' */
|
||||
|
||||
+94
-73
@@ -34,7 +34,8 @@ typedef struct {
|
||||
/* Initialize the plasma manager. This function initializes the event loop
|
||||
* of the plasma manager, and stores the address 'store_socket_name' of
|
||||
* the local plasma store socket. */
|
||||
void init_plasma_manager(plasma_manager_state *s, const char* store_socket_name) {
|
||||
void init_plasma_manager(plasma_manager_state* s,
|
||||
const char* store_socket_name) {
|
||||
s->loop = malloc(sizeof(event_loop));
|
||||
event_loop_init(s->loop);
|
||||
s->store_socket_name = store_socket_name;
|
||||
@@ -45,36 +46,47 @@ void init_plasma_manager(plasma_manager_state *s, const char* store_socket_name)
|
||||
* the data header to the other object manager. */
|
||||
void initiate_transfer(plasma_manager_state* s, plasma_request* req) {
|
||||
int store_conn = plasma_store_connect(s->store_socket_name);
|
||||
plasma_buffer buf = { .object_id = req->object_id, .writable = 0 };
|
||||
plasma_buffer buf = {.object_id = req->object_id, .writable = 0};
|
||||
plasma_get(store_conn, req->object_id, &buf.size, &buf.data);
|
||||
|
||||
|
||||
char ip_addr[16];
|
||||
snprintf(ip_addr, 32, "%d.%d.%d.%d",
|
||||
req->addr[0], req->addr[1],
|
||||
req->addr[2], req->addr[3]);
|
||||
snprintf(ip_addr, 32, "%d.%d.%d.%d", req->addr[0], req->addr[1], req->addr[2],
|
||||
req->addr[3]);
|
||||
|
||||
int fd = plasma_manager_connect(&ip_addr[0], req->port);
|
||||
data_connection conn = { .type = DATA_CONNECTION_WRITE, .store_conn = store_conn, .buf = buf, .cursor = 0 };
|
||||
data_connection conn = {.type = DATA_CONNECTION_WRITE,
|
||||
.store_conn = store_conn,
|
||||
.buf = buf,
|
||||
.cursor = 0};
|
||||
event_loop_attach(s->loop, CONNECTION_DATA, &conn, fd, POLLOUT);
|
||||
|
||||
plasma_request manager_req = { .type = PLASMA_DATA, .object_id = req->object_id, .size = buf.size };
|
||||
plasma_request manager_req = {
|
||||
.type = PLASMA_DATA, .object_id = req->object_id, .size = buf.size};
|
||||
plasma_send(fd, &manager_req);
|
||||
}
|
||||
|
||||
/* Start reading data from another object manager.
|
||||
* Initializes the object we are going to write to in the
|
||||
* local plasma store and then switches the data socket to reading mode. */
|
||||
void start_reading_data(int64_t index, plasma_manager_state* s, plasma_request* req) {
|
||||
void start_reading_data(int64_t index,
|
||||
plasma_manager_state* s,
|
||||
plasma_request* req) {
|
||||
int store_conn = plasma_store_connect(s->store_socket_name);
|
||||
plasma_buffer buf = { .object_id = req->object_id, .size = req->size, .writable = 1 };
|
||||
plasma_buffer buf = {
|
||||
.object_id = req->object_id, .size = req->size, .writable = 1};
|
||||
plasma_create(store_conn, req->object_id, req->size, &buf.data);
|
||||
data_connection conn = { .type = DATA_CONNECTION_READ, .store_conn = store_conn, .buf = buf, .cursor = 0 };
|
||||
data_connection conn = {.type = DATA_CONNECTION_READ,
|
||||
.store_conn = store_conn,
|
||||
.buf = buf,
|
||||
.cursor = 0};
|
||||
event_loop_set_connection(s->loop, index, &conn);
|
||||
}
|
||||
|
||||
/* Handle a command request that came in through a socket (transfering data,
|
||||
* or accepting incoming data). */
|
||||
void process_command(int64_t id, plasma_manager_state* state, plasma_request* req) {
|
||||
void process_command(int64_t id,
|
||||
plasma_manager_state* state,
|
||||
plasma_request* req) {
|
||||
switch (req->type) {
|
||||
case PLASMA_TRANSFER:
|
||||
LOG_INFO("transfering object to manager with port %d", req->port);
|
||||
@@ -91,63 +103,66 @@ void process_command(int64_t id, plasma_manager_state* state, plasma_request* re
|
||||
}
|
||||
|
||||
/* Handle data or command event incoming on socket with index "index". */
|
||||
void read_from_socket(plasma_manager_state* state, struct pollfd *waiting, int64_t index, plasma_request* req) {
|
||||
void read_from_socket(plasma_manager_state* state,
|
||||
struct pollfd* waiting,
|
||||
int64_t index,
|
||||
plasma_request* req) {
|
||||
ssize_t r, s;
|
||||
data_connection *conn = event_loop_get_connection(state->loop, index);
|
||||
data_connection* conn = event_loop_get_connection(state->loop, index);
|
||||
switch (conn->type) {
|
||||
case DATA_CONNECTION_HEADER:
|
||||
r = read(waiting->fd, req, sizeof(plasma_request));
|
||||
if (r == -1) {
|
||||
LOG_ERR("read error");
|
||||
} else if (r == 0) {
|
||||
LOG_INFO("connection with id %" PRId64 " disconnected", index);
|
||||
event_loop_detach(state->loop, index, 1);
|
||||
case DATA_CONNECTION_HEADER:
|
||||
r = read(waiting->fd, req, sizeof(plasma_request));
|
||||
if (r == -1) {
|
||||
LOG_ERR("read error");
|
||||
} else if (r == 0) {
|
||||
LOG_INFO("connection with id %" PRId64 " disconnected", index);
|
||||
event_loop_detach(state->loop, index, 1);
|
||||
} else {
|
||||
process_command(index, state, req);
|
||||
}
|
||||
break;
|
||||
case DATA_CONNECTION_READ:
|
||||
LOG_DEBUG("polled DATA_CONNECTION_READ");
|
||||
r = read(waiting->fd, conn->buf.data + conn->cursor, BUFSIZE);
|
||||
if (r == -1) {
|
||||
LOG_ERR("read error");
|
||||
} else if (r == 0) {
|
||||
LOG_INFO("end of file");
|
||||
} else {
|
||||
conn->cursor += r;
|
||||
}
|
||||
if (r == 0) {
|
||||
LOG_DEBUG("reading on channel %" PRId64 " finished", index);
|
||||
plasma_seal(conn->store_conn, conn->buf.object_id);
|
||||
close(conn->store_conn);
|
||||
event_loop_detach(state->loop, index, 1);
|
||||
}
|
||||
break;
|
||||
case DATA_CONNECTION_WRITE:
|
||||
LOG_DEBUG("polled DATA_CONNECTION_WRITE");
|
||||
s = conn->buf.size - conn->cursor;
|
||||
if (s > BUFSIZE)
|
||||
s = BUFSIZE;
|
||||
r = write(waiting->fd, conn->buf.data + conn->cursor, s);
|
||||
if (r != s) {
|
||||
if (r > 0) {
|
||||
LOG_ERR("partial write on fd %d", waiting->fd);
|
||||
} else {
|
||||
process_command(index, state, req);
|
||||
LOG_ERR("write error");
|
||||
exit(-1);
|
||||
}
|
||||
break;
|
||||
case DATA_CONNECTION_READ:
|
||||
LOG_DEBUG("polled DATA_CONNECTION_READ");
|
||||
r = read(waiting->fd, conn->buf.data + conn->cursor, BUFSIZE);
|
||||
if (r == -1) {
|
||||
LOG_ERR("read error");
|
||||
} else if (r == 0) {
|
||||
LOG_INFO("end of file");
|
||||
} else {
|
||||
conn->cursor += r;
|
||||
}
|
||||
if (r == 0) {
|
||||
LOG_DEBUG("reading on channel %" PRId64 " finished", index);
|
||||
plasma_seal(conn->store_conn, conn->buf.object_id);
|
||||
close(conn->store_conn);
|
||||
event_loop_detach(state->loop, index, 1);
|
||||
}
|
||||
break;
|
||||
case DATA_CONNECTION_WRITE:
|
||||
LOG_DEBUG("polled DATA_CONNECTION_WRITE");
|
||||
s = conn->buf.size - conn->cursor;
|
||||
if (s > BUFSIZE)
|
||||
s = BUFSIZE;
|
||||
r = write(waiting->fd, conn->buf.data + conn->cursor, s);
|
||||
if (r != s) {
|
||||
if (r > 0) {
|
||||
LOG_ERR("partial write on fd %d", waiting->fd);
|
||||
} else {
|
||||
LOG_ERR("write error");
|
||||
exit(-1);
|
||||
}
|
||||
} else {
|
||||
conn->cursor += r;
|
||||
}
|
||||
if (r == 0) {
|
||||
LOG_DEBUG("writing on channel %" PRId64 " finished", index);
|
||||
close(conn->store_conn);
|
||||
event_loop_detach(state->loop, index, 1);
|
||||
}
|
||||
break;
|
||||
default:
|
||||
LOG_ERR("invalid connection type");
|
||||
exit(-1);
|
||||
} else {
|
||||
conn->cursor += r;
|
||||
}
|
||||
if (r == 0) {
|
||||
LOG_DEBUG("writing on channel %" PRId64 " finished", index);
|
||||
close(conn->store_conn);
|
||||
event_loop_detach(state->loop, index, 1);
|
||||
}
|
||||
break;
|
||||
default:
|
||||
LOG_ERR("invalid connection type");
|
||||
exit(-1);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -163,7 +178,7 @@ void run_event_loop(int sock, plasma_manager_state* s) {
|
||||
exit(-1);
|
||||
}
|
||||
for (int i = 0; i < event_loop_size(s->loop); ++i) {
|
||||
struct pollfd *waiting = event_loop_get(s->loop, i);
|
||||
struct pollfd* waiting = event_loop_get(s->loop, i);
|
||||
if (waiting->revents == 0)
|
||||
continue;
|
||||
if (waiting->fd == sock) {
|
||||
@@ -176,7 +191,7 @@ void run_event_loop(int sock, plasma_manager_state* s) {
|
||||
}
|
||||
break;
|
||||
}
|
||||
data_connection conn = { .type = DATA_CONNECTION_HEADER };
|
||||
data_connection conn = {.type = DATA_CONNECTION_HEADER};
|
||||
event_loop_attach(s->loop, CONNECTION_DATA, &conn, new_socket, POLLIN);
|
||||
LOG_INFO("new connection with id %" PRId64, event_loop_size(s->loop));
|
||||
} else {
|
||||
@@ -186,7 +201,9 @@ void run_event_loop(int sock, plasma_manager_state* s) {
|
||||
}
|
||||
}
|
||||
|
||||
void start_server(const char *store_socket_name, const char* master_addr, int port) {
|
||||
void start_server(const char* store_socket_name,
|
||||
const char* master_addr,
|
||||
int port) {
|
||||
struct sockaddr_in name;
|
||||
int sock = socket(PF_INET, SOCK_STREAM, 0);
|
||||
if (sock < 0) {
|
||||
@@ -203,7 +220,7 @@ void start_server(const char *store_socket_name, const char* master_addr, int po
|
||||
close(sock);
|
||||
exit(-1);
|
||||
}
|
||||
setsockopt(sock, SOL_SOCKET, SO_REUSEADDR, &on, sizeof (on));
|
||||
setsockopt(sock, SOL_SOCKET, SO_REUSEADDR, &on, sizeof(on));
|
||||
if (bind(sock, (struct sockaddr*) &name, sizeof(name)) < 0) {
|
||||
LOG_ERR("could not bind socket");
|
||||
exit(-1);
|
||||
@@ -220,9 +237,9 @@ void start_server(const char *store_socket_name, const char* master_addr, int po
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
/* Socket name of the plasma store this manager is connected to. */
|
||||
char *store_socket_name = NULL;
|
||||
char* store_socket_name = NULL;
|
||||
/* IP address of this node. */
|
||||
char *master_addr = NULL;
|
||||
char* master_addr = NULL;
|
||||
/* Port number the manager should use. */
|
||||
int port;
|
||||
int c;
|
||||
@@ -243,11 +260,15 @@ int main(int argc, char* argv[]) {
|
||||
}
|
||||
}
|
||||
if (!store_socket_name) {
|
||||
LOG_ERR("please specify socket for connecting to the plasma store with -s switch");
|
||||
LOG_ERR(
|
||||
"please specify socket for connecting to the plasma store with -s "
|
||||
"switch");
|
||||
exit(-1);
|
||||
}
|
||||
if (!master_addr) {
|
||||
LOG_ERR("please specify ip address of the current host in the format 123.456.789.10 with -m switch");
|
||||
LOG_ERR(
|
||||
"please specify ip address of the current host in the format "
|
||||
"123.456.789.10 with -m switch");
|
||||
exit(-1);
|
||||
}
|
||||
start_server(store_socket_name, master_addr, port);
|
||||
|
||||
@@ -7,11 +7,7 @@
|
||||
/* The buffer size in bytes. Data will get transfered in multiples of this */
|
||||
#define BUFSIZE 4096
|
||||
|
||||
enum connection_type {
|
||||
CONNECTION_REDIS,
|
||||
CONNECTION_LISTENER,
|
||||
CONNECTION_DATA
|
||||
};
|
||||
enum connection_type { CONNECTION_REDIS, CONNECTION_LISTENER, CONNECTION_DATA };
|
||||
|
||||
enum data_connection_type {
|
||||
/* Connection to send commands and metadata to the manager. */
|
||||
|
||||
+19
-18
@@ -9,7 +9,6 @@
|
||||
* It keeps a hash table that maps object_ids (which are 20 byte long,
|
||||
* just enough to store and SHA1 hash) to memory mapped files. */
|
||||
|
||||
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <unistd.h>
|
||||
@@ -30,7 +29,7 @@
|
||||
|
||||
typedef struct {
|
||||
/* Event loop for the plasma store. */
|
||||
event_loop *loop;
|
||||
event_loop* loop;
|
||||
} plasma_store_state;
|
||||
|
||||
void init_state(plasma_store_state* s) {
|
||||
@@ -103,41 +102,42 @@ void create_object(int conn, plasma_request* req) {
|
||||
LOG_ERR("could not create shared memory buffer");
|
||||
exit(-1);
|
||||
}
|
||||
object_table_entry *entry = malloc(sizeof(object_table_entry));
|
||||
object_table_entry* entry = malloc(sizeof(object_table_entry));
|
||||
memcpy(&entry->object_id, &req->object_id, 20);
|
||||
entry->info.size = req->size;
|
||||
/* TODO(pcm): set the other fields */
|
||||
entry->fd = fd;
|
||||
HASH_ADD(handle, open_objects, object_id, sizeof(plasma_id), entry);
|
||||
plasma_reply reply = { PLASMA_OBJECT, req->size };
|
||||
plasma_reply reply = {PLASMA_OBJECT, req->size};
|
||||
send_fd(conn, fd, (char*) &reply, sizeof(plasma_reply));
|
||||
}
|
||||
|
||||
/* Get an object from the hash table. */
|
||||
void get_object(int conn, plasma_request* req) {
|
||||
object_table_entry *entry;
|
||||
object_table_entry* entry;
|
||||
HASH_FIND(handle, sealed_objects, &req->object_id, sizeof(plasma_id), entry);
|
||||
if (entry) {
|
||||
plasma_reply reply = { PLASMA_OBJECT, entry->info.size };
|
||||
plasma_reply reply = {PLASMA_OBJECT, entry->info.size};
|
||||
send_fd(conn, entry->fd, (char*) &reply, sizeof(plasma_reply));
|
||||
} else {
|
||||
LOG_INFO("object not in hash table of sealed objects");
|
||||
int fd[2];
|
||||
socketpair(AF_UNIX, SOCK_STREAM, 0, fd);
|
||||
object_notify_entry *notify_entry = malloc(sizeof(object_notify_entry));
|
||||
object_notify_entry* notify_entry = malloc(sizeof(object_notify_entry));
|
||||
memcpy(¬ify_entry->object_id, &req->object_id, 20);
|
||||
notify_entry->conn[notify_entry->num_waiting] = fd[0];
|
||||
notify_entry->num_waiting += 1;
|
||||
HASH_ADD(handle, objects_notify, object_id, sizeof(plasma_id), notify_entry);
|
||||
plasma_reply reply = { PLASMA_FUTURE, -1 };
|
||||
HASH_ADD(handle, objects_notify, object_id, sizeof(plasma_id),
|
||||
notify_entry);
|
||||
plasma_reply reply = {PLASMA_FUTURE, -1};
|
||||
send_fd(conn, fd[1], (char*) &reply, sizeof(plasma_reply));
|
||||
}
|
||||
}
|
||||
|
||||
/* Seal an object that has been created in the hash table. */
|
||||
void seal_object(int conn, plasma_request* req) {
|
||||
LOG_INFO("sealing object"); // TODO(pcm): add object_id here
|
||||
object_table_entry *entry;
|
||||
LOG_INFO("sealing object"); // TODO(pcm): add object_id here
|
||||
object_table_entry* entry;
|
||||
HASH_FIND(handle, open_objects, &req->object_id, sizeof(plasma_id), entry);
|
||||
if (!entry) {
|
||||
return; /* TODO(pcm): return error */
|
||||
@@ -148,11 +148,12 @@ void seal_object(int conn, plasma_request* req) {
|
||||
HASH_ADD(handle, sealed_objects, object_id, sizeof(plasma_id), entry);
|
||||
/* Inform processes that the object is ready now. */
|
||||
object_notify_entry* notify_entry;
|
||||
HASH_FIND(handle, objects_notify, &req->object_id, sizeof(plasma_id), notify_entry);
|
||||
HASH_FIND(handle, objects_notify, &req->object_id, sizeof(plasma_id),
|
||||
notify_entry);
|
||||
if (!notify_entry) {
|
||||
return;
|
||||
}
|
||||
plasma_reply reply = { PLASMA_OBJECT, size };
|
||||
plasma_reply reply = {PLASMA_OBJECT, size};
|
||||
for (int i = 0; i < notify_entry->num_waiting; ++i) {
|
||||
send_fd(notify_entry->conn[i], fd, (char*) &reply, sizeof(plasma_reply));
|
||||
close(notify_entry->conn[i]);
|
||||
@@ -190,7 +191,7 @@ void run_event_loop(int socket) {
|
||||
exit(-1);
|
||||
}
|
||||
for (int i = 0; i < event_loop_size(state.loop); ++i) {
|
||||
struct pollfd *waiting = event_loop_get(state.loop, i);
|
||||
struct pollfd* waiting = event_loop_get(state.loop, i);
|
||||
if (waiting->revents == 0)
|
||||
continue;
|
||||
if (waiting->fd == socket) {
|
||||
@@ -230,7 +231,7 @@ void start_server(char* socket_name) {
|
||||
exit(-1);
|
||||
}
|
||||
int on = 1;
|
||||
if (setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, (char*)&on, sizeof(on)) < 0) {
|
||||
if (setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, (char*) &on, sizeof(on)) < 0) {
|
||||
LOG_ERR("setsockopt failed");
|
||||
close(fd);
|
||||
exit(-1);
|
||||
@@ -244,15 +245,15 @@ void start_server(char* socket_name) {
|
||||
struct sockaddr_un addr;
|
||||
memset(&addr, 0, sizeof(addr));
|
||||
addr.sun_family = AF_UNIX;
|
||||
strncpy(addr.sun_path, socket_name, sizeof(addr.sun_path)-1);
|
||||
strncpy(addr.sun_path, socket_name, sizeof(addr.sun_path) - 1);
|
||||
unlink(socket_name);
|
||||
bind(fd, (struct sockaddr*)&addr, sizeof(addr));
|
||||
bind(fd, (struct sockaddr*) &addr, sizeof(addr));
|
||||
listen(fd, 5);
|
||||
run_event_loop(fd);
|
||||
}
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
char *socket_name = NULL;
|
||||
char* socket_name = NULL;
|
||||
int c;
|
||||
while ((c = getopt(argc, argv, "s:")) != -1) {
|
||||
switch (c) {
|
||||
|
||||
Reference in New Issue
Block a user