diff --git a/src/plasma/plasma_manager.cc b/src/plasma/plasma_manager.cc index 4aaad9663..40904f251 100644 --- a/src/plasma/plasma_manager.cc +++ b/src/plasma/plasma_manager.cc @@ -104,7 +104,7 @@ typedef struct { ObjectID object_id; /** Handle for the uthash table. */ UT_hash_handle hh; -} available_object; +} AvailableObject; typedef struct { /** The ID of the object we are fetching or waiting for. */ @@ -121,27 +121,27 @@ typedef struct { /** Handle for the uthash table in the manager state that keeps track of * outstanding fetch requests. */ UT_hash_handle hh; -} fetch_request; +} FetchRequest; /** * There are fundamentally two data structures used for handling wait requests. - * There is the "wait_request" struct and the "object_wait_requests" struct. A - * wait_request keeps track of all of the object IDs that a wait_request is - * waiting for. An object_wait_requests struct keeps track of all of the - * wait_request structs that are waiting for a particular object iD. The - * plasma_manager_state contains a hash table mapping object IDs to their - * coresponding object_wait_requests structs. + * There is the "wait_request" struct and the "ObjectWaitRequests" struct. A + * WaitRequest keeps track of all of the object IDs that a WaitRequest is + * waiting for. An ObjectWaitRequests struct keeps track of all of the + * WaitRequest structs that are waiting for a particular object iD. The + * PlasmaManagerState contains a hash table mapping object IDs to their + * coresponding ObjectWaitRequests structs. * * These data structures are updated by several methods: - * - add_wait_request_for_object adds a wait_request to the - * object_wait_requests struct corresponding to a particular object ID. This + * - add_wait_request_for_object adds a WaitRequest to the + * ObjectWaitRequests struct corresponding to a particular object ID. This * is called when a client calls plasma_wait. - * - remove_wait_request_for_object removes a wait_request from an - * object_wait_requests struct. When a wait request returns, this method is - * called for all of the object IDs involved in that wait_request. - * - update_object_wait_requests removes an object_wait_requests struct and - * does some processing for each wait_request involved in that - * object_wait_requests struct. + * - remove_wait_request_for_object removes a WaitRequest from an + * ObjectWaitRequests struct. When a wait request returns, this method is + * called for all of the object IDs involved in that WaitRequest. + * - update_object_wait_requests removes an ObjectWaitRequests struct and + * does some processing for each WaitRequest involved in that + * ObjectWaitRequests struct. */ typedef struct { /** The client connection that called wait. */ @@ -160,11 +160,11 @@ typedef struct { /** The number of object requests in this wait request that are already * satisfied. */ int64_t num_satisfied; -} wait_request; +} WaitRequest; /** This is used to define the utarray of wait requests in the - * object_wait_requests struct. */ -UT_icd wait_request_icd = {sizeof(wait_request *), NULL, NULL, NULL}; + * ObjectWaitRequests struct. */ +UT_icd wait_request_icd = {sizeof(WaitRequest *), NULL, NULL, NULL}; typedef struct { /** The ID of the object. This is used as a key in a hash table. */ @@ -174,9 +174,9 @@ typedef struct { /** Handle for the uthash table in the manager state that keeps track of the * wait requests involving this object ID. */ UT_hash_handle hh; -} object_wait_requests; +} ObjectWaitRequests; -struct plasma_manager_state { +struct PlasmaManagerState { /** Event loop. */ event_loop *loop; /** Connection to the local plasma store for reading or writing data. */ @@ -192,27 +192,27 @@ struct plasma_manager_state { int port; /** Hash table of outstanding fetch requests. The key is the object ID. The * value is the data needed to perform the fetch. */ - fetch_request *fetch_requests; + FetchRequest *fetch_requests; /** A hash table mapping object IDs to a vector of the wait requests that * are waiting for the object to arrive locally. */ - object_wait_requests *object_wait_requests_local; + ObjectWaitRequests *object_wait_requests_local; /** A hash table mapping object IDs to a vector of the wait requests that * are waiting for the object to be available somewhere in the system. */ - object_wait_requests *object_wait_requests_remote; + ObjectWaitRequests *object_wait_requests_remote; /** Initialize an empty hash map for the cache of local available object. */ - available_object *local_available_objects; + AvailableObject *local_available_objects; /** Buffer that holds memory for serializing plasma protocol messages. */ protocol_builder *builder; }; -plasma_manager_state *g_manager_state = NULL; +PlasmaManagerState *g_manager_state = NULL; /* The context for fetch and wait requests. These are per client, per object. */ -struct client_object_request { +struct ClientObjectRequest { /** The ID of the object we are fetching or waiting for. */ ObjectID object_id; - /** The client connection context, shared between other - * client_object_requests for the same client. */ + /** The client connection context, shared between other ClientObjectRequest + * structs for the same client. */ ClientConnection *client_conn; /** The ID for the timer that will time out the current request to the state * database or another plasma manager. */ @@ -221,7 +221,7 @@ struct client_object_request { * timeout. */ int num_retries; /** Handle for a linked list. */ - client_object_request *next; + ClientObjectRequest *next; /** Pointer to the array containing the manager locations of * this object. */ char **manager_vector; @@ -244,16 +244,20 @@ struct client_object_request { struct ClientConnection { /** Current state for this plasma manager. This is shared * between all client connections to the plasma manager. */ - plasma_manager_state *manager_state; + PlasmaManagerState *manager_state; /** Current position in the buffer. */ int64_t cursor; /** Buffer that this connection is reading from. If this is a connection to * write data to another plasma store, then it is a linked * list of buffers to write. */ /* TODO(swang): Split into two queues, data transfers and data requests. */ - plasma_request_buffer *transfer_queue; + PlasmaRequestBuffer *transfer_queue; + /* A set of object IDs which are queued in the transfer_queue and waiting to + * be sent. This is used to avoid sending the same object ID to the same + * manager multiple times. */ + PlasmaRequestBuffer *pending_object_transfers; /** Buffer used to receive transfers (data fetches) we want to ignore */ - plasma_request_buffer *ignore_buffer; + PlasmaRequestBuffer *ignore_buffer; /** File descriptor for the socket connected to the other * plasma manager. */ int fd; @@ -261,7 +265,7 @@ struct ClientConnection { int64_t timer_id; /** The objects that we are waiting for and their callback * contexts, for either a fetch or a wait operation. */ - client_object_request *active_objects; + ClientObjectRequest *active_objects; /** The number of objects that we have left to return for * this fetch or wait operation. */ int num_return_objects; @@ -273,8 +277,8 @@ struct ClientConnection { UT_hash_handle manager_hh; }; -object_wait_requests **object_wait_requests_table_ptr_from_type( - plasma_manager_state *manager_state, +ObjectWaitRequests **object_wait_requests_table_ptr_from_type( + PlasmaManagerState *manager_state, int type) { /* We use different types of hash tables for different requests. */ if (type == PLASMA_QUERY_LOCAL) { @@ -286,21 +290,21 @@ object_wait_requests **object_wait_requests_table_ptr_from_type( } } -void add_wait_request_for_object(plasma_manager_state *manager_state, +void add_wait_request_for_object(PlasmaManagerState *manager_state, ObjectID object_id, int type, - wait_request *wait_req) { - object_wait_requests **object_wait_requests_table_ptr = + WaitRequest *wait_req) { + ObjectWaitRequests **object_wait_requests_table_ptr = object_wait_requests_table_ptr_from_type(manager_state, type); - object_wait_requests *object_wait_reqs; + ObjectWaitRequests *object_wait_reqs; HASH_FIND(hh, *object_wait_requests_table_ptr, &object_id, sizeof(object_id), object_wait_reqs); /* If there are currently no wait requests involving this object ID, create a - * new object_wait_requests struct for this object ID and add it to the hash + * new ObjectWaitRequests struct for this object ID and add it to the hash * table. */ if (object_wait_reqs == NULL) { object_wait_reqs = - (object_wait_requests *) malloc(sizeof(object_wait_requests)); + (ObjectWaitRequests *) malloc(sizeof(ObjectWaitRequests)); object_wait_reqs->object_id = object_id; utarray_new(object_wait_reqs->wait_requests, &wait_request_icd); HASH_ADD(hh, *object_wait_requests_table_ptr, object_id, @@ -311,13 +315,13 @@ void add_wait_request_for_object(plasma_manager_state *manager_state, utarray_push_back(object_wait_reqs->wait_requests, &wait_req); } -void remove_wait_request_for_object(plasma_manager_state *manager_state, +void remove_wait_request_for_object(PlasmaManagerState *manager_state, ObjectID object_id, int type, - wait_request *wait_req) { - object_wait_requests **object_wait_requests_table_ptr = + WaitRequest *wait_req) { + ObjectWaitRequests **object_wait_requests_table_ptr = object_wait_requests_table_ptr_from_type(manager_state, type); - object_wait_requests *object_wait_reqs; + ObjectWaitRequests *object_wait_reqs; HASH_FIND(hh, *object_wait_requests_table_ptr, &object_id, sizeof(object_id), object_wait_reqs); /* If there is a vector of wait requests for this object ID, and if this @@ -325,8 +329,8 @@ void remove_wait_request_for_object(plasma_manager_state *manager_state, * vector. */ if (object_wait_reqs != NULL) { for (int i = 0; i < utarray_len(object_wait_reqs->wait_requests); ++i) { - wait_request **wait_req_ptr = - (wait_request **) utarray_eltptr(object_wait_reqs->wait_requests, i); + WaitRequest **wait_req_ptr = + (WaitRequest **) utarray_eltptr(object_wait_reqs->wait_requests, i); if (*wait_req_ptr == wait_req) { /* Remove the wait request from the array. */ utarray_erase(object_wait_reqs->wait_requests, i, 1); @@ -339,8 +343,8 @@ void remove_wait_request_for_object(plasma_manager_state *manager_state, } } -void remove_wait_request(plasma_manager_state *manager_state, - wait_request *wait_req) { +void remove_wait_request(PlasmaManagerState *manager_state, + WaitRequest *wait_req) { if (wait_req->timer != -1) { CHECK(event_loop_remove_timer(manager_state->loop, wait_req->timer) == AE_OK); @@ -349,8 +353,8 @@ void remove_wait_request(plasma_manager_state *manager_state, free(wait_req); } -void return_from_wait(plasma_manager_state *manager_state, - wait_request *wait_req) { +void return_from_wait(PlasmaManagerState *manager_state, + WaitRequest *wait_req) { /* Send the reply to the client. */ warn_if_sigpipe(plasma_send_WaitReply( wait_req->client_conn->fd, manager_state->builder, @@ -367,14 +371,14 @@ void return_from_wait(plasma_manager_state *manager_state, remove_wait_request(manager_state, wait_req); } -void update_object_wait_requests(plasma_manager_state *manager_state, +void update_object_wait_requests(PlasmaManagerState *manager_state, ObjectID obj_id, int type, int status) { - object_wait_requests **object_wait_requests_table_ptr = + ObjectWaitRequests **object_wait_requests_table_ptr = object_wait_requests_table_ptr_from_type(manager_state, type); /* Update the in-progress wait requests in the specified table. */ - object_wait_requests *object_wait_reqs; + ObjectWaitRequests *object_wait_reqs; HASH_FIND(hh, *object_wait_requests_table_ptr, &obj_id, sizeof(obj_id), object_wait_reqs); if (object_wait_reqs != NULL) { @@ -387,9 +391,9 @@ void update_object_wait_requests(plasma_manager_state *manager_state, * are removed from the array. */ int index = 0; for (int i = 0; i < num_requests; ++i) { - wait_request **wait_req_ptr = (wait_request **) utarray_eltptr( + WaitRequest **wait_req_ptr = (WaitRequest **) utarray_eltptr( object_wait_reqs->wait_requests, index); - wait_request *wait_req = *wait_req_ptr; + WaitRequest *wait_req = *wait_req_ptr; wait_req->num_satisfied += 1; /* Mark the object as present in the wait request. */ int j = 0; @@ -422,17 +426,17 @@ void update_object_wait_requests(plasma_manager_state *manager_state, } } -fetch_request *create_fetch_request(plasma_manager_state *manager_state, - ObjectID object_id) { - fetch_request *fetch_req = (fetch_request *) malloc(sizeof(fetch_request)); +FetchRequest *create_fetch_request(PlasmaManagerState *manager_state, + ObjectID object_id) { + FetchRequest *fetch_req = (FetchRequest *) malloc(sizeof(FetchRequest)); fetch_req->object_id = object_id; fetch_req->manager_count = 0; fetch_req->manager_vector = NULL; return fetch_req; } -void remove_fetch_request(plasma_manager_state *manager_state, - fetch_request *fetch_req) { +void remove_fetch_request(PlasmaManagerState *manager_state, + FetchRequest *fetch_req) { /* Remove the fetch request from the table of fetch requests. */ HASH_DELETE(hh, manager_state->fetch_requests, fetch_req); /* Free the fetch request and everything in it. */ @@ -445,14 +449,14 @@ void remove_fetch_request(plasma_manager_state *manager_state, free(fetch_req); } -plasma_manager_state *init_plasma_manager_state(const char *store_socket_name, - const char *manager_socket_name, - const char *manager_addr, - int manager_port, - const char *db_addr, - int db_port) { - plasma_manager_state *state = - (plasma_manager_state *) malloc(sizeof(plasma_manager_state)); +PlasmaManagerState *PlasmaManagerState_init(const char *store_socket_name, + const char *manager_socket_name, + const char *manager_addr, + int manager_port, + const char *db_addr, + int db_port) { + PlasmaManagerState *state = + (PlasmaManagerState *) malloc(sizeof(PlasmaManagerState)); state->loop = event_loop_create(); state->plasma_conn = plasma_connect(store_socket_name, NULL, PLASMA_DEFAULT_RELEASE_DELAY); @@ -497,23 +501,35 @@ plasma_manager_state *init_plasma_manager_state(const char *store_socket_name, return state; } -void destroy_plasma_manager_state(plasma_manager_state *state) { +void PlasmaManagerState_free(PlasmaManagerState *state) { ClientConnection *manager_conn, *tmp; HASH_ITER(manager_hh, state->manager_connections, manager_conn, tmp) { HASH_DELETE(manager_hh, state->manager_connections, manager_conn); - plasma_request_buffer *head = manager_conn->transfer_queue; + + /* Free the hash table of object IDs that are waiting to be transferred. */ + PlasmaRequestBuffer *request_buffer, *tmp_buffer; + HASH_ITER(hh, manager_conn->pending_object_transfers, request_buffer, + tmp_buffer) { + /* We do not free the PlasmaRequestBuffer here because it is also in the + * transfer queue and will be freed below. */ + HASH_DELETE(hh, manager_conn->pending_object_transfers, request_buffer); + } + + /* Free the transfer queue. */ + PlasmaRequestBuffer *head = manager_conn->transfer_queue; while (head) { - LL_DELETE(manager_conn->transfer_queue, head); + DL_DELETE(manager_conn->transfer_queue, head); free(head); head = manager_conn->transfer_queue; } + /* Close the manager connection and free the remaining state. */ close(manager_conn->fd); free(manager_conn->ip_addr_port); free(manager_conn); } if (state->fetch_requests != NULL) { - fetch_request *fetch_req, *tmp; + FetchRequest *fetch_req, *tmp; HASH_ITER(hh, state->fetch_requests, fetch_req, tmp) { remove_fetch_request(state, fetch_req); } @@ -525,7 +541,7 @@ void destroy_plasma_manager_state(plasma_manager_state *state) { free(state); } -event_loop *get_event_loop(plasma_manager_state *state) { +event_loop *get_event_loop(PlasmaManagerState *state) { return state->loop; } @@ -536,7 +552,7 @@ void process_message(event_loop *loop, void *context, int events); -void write_object_chunk(ClientConnection *conn, plasma_request_buffer *buf) { +void write_object_chunk(ClientConnection *conn, PlasmaRequestBuffer *buf) { LOG_DEBUG("Writing data to fd %d", conn->fd); ssize_t r, s; /* Try to write one BUFSIZE at a time. */ @@ -571,7 +587,7 @@ void send_queued_request(event_loop *loop, void *context, int events) { ClientConnection *conn = (ClientConnection *) context; - plasma_manager_state *state = conn->manager_state; + PlasmaManagerState *state = conn->manager_state; if (conn->transfer_queue == NULL) { /* If there are no objects to transfer, temporarily remove this connection @@ -581,7 +597,7 @@ void send_queued_request(event_loop *loop, return; } - plasma_request_buffer *buf = conn->transfer_queue; + PlasmaRequestBuffer *buf = conn->transfer_queue; switch (buf->type) { case MessageType_PlasmaDataRequest: warn_if_sigpipe( @@ -605,14 +621,19 @@ void send_queued_request(event_loop *loop, LOG_FATAL("Buffered request has unknown type."); } - /* We are done sending this request. */ + /* If we are done sending this request, remove it from the transfer queue. */ if (conn->cursor == 0) { - LL_DELETE(conn->transfer_queue, buf); + if (buf->type == MessageType_PlasmaDataReply) { + /* If we just finished sending an object to a remote manager, then remove + * the object from the hash table of pending transfer requests. */ + HASH_DELETE(hh, conn->pending_object_transfers, buf); + } + DL_DELETE(conn->transfer_queue, buf); free(buf); } } -int read_object_chunk(ClientConnection *conn, plasma_request_buffer *buf) { +int read_object_chunk(ClientConnection *conn, PlasmaRequestBuffer *buf) { LOG_DEBUG("Reading data from fd %d to %p", conn->fd, buf->data + conn->cursor); ssize_t r, s; @@ -647,7 +668,7 @@ void process_data_chunk(event_loop *loop, int events) { /* Read the object chunk. */ ClientConnection *conn = (ClientConnection *) context; - plasma_request_buffer *buf = conn->transfer_queue; + PlasmaRequestBuffer *buf = conn->transfer_queue; int done = read_object_chunk(conn, buf); if (!done) { return; @@ -661,7 +682,7 @@ void process_data_chunk(event_loop *loop, plasma_seal(conn->manager_state->plasma_conn, buf->object_id); plasma_release(conn->manager_state->plasma_conn, buf->object_id); /* Remove the request buffer used for reading this object's data. */ - LL_DELETE(conn->transfer_queue, buf); + DL_DELETE(conn->transfer_queue, buf); free(buf); /* Switch to listening for requests from this socket, instead of reading * object data. */ @@ -675,7 +696,7 @@ void ignore_data_chunk(event_loop *loop, int events) { /* Read the object chunk. */ ClientConnection *conn = (ClientConnection *) context; - plasma_request_buffer *buf = conn->ignore_buffer; + PlasmaRequestBuffer *buf = conn->ignore_buffer; /* Just read the transferred data into ignore_buf and then drop (free) it. */ int done = read_object_chunk(conn, buf); @@ -691,7 +712,7 @@ void ignore_data_chunk(event_loop *loop, event_loop_add_file(loop, data_sock, EVENT_LOOP_READ, process_message, conn); } -ClientConnection *get_manager_connection(plasma_manager_state *state, +ClientConnection *get_manager_connection(PlasmaManagerState *state, const char *ip_addr, int port) { /* TODO(swang): Should probably check whether ip_addr and port belong to us. @@ -712,6 +733,7 @@ ClientConnection *get_manager_connection(plasma_manager_state *state, manager_conn->fd = fd; manager_conn->manager_state = state; manager_conn->transfer_queue = NULL; + manager_conn->pending_object_transfers = NULL; manager_conn->cursor = 0; manager_conn->ip_addr_port = strdup(utstring_body(ip_addr_port)); HASH_ADD_KEYPTR(manager_hh, @@ -733,12 +755,11 @@ void process_transfer_request(event_loop *loop, /* If there is already a request in the transfer queue with the same object * ID, do not add the transfer request. */ - plasma_request_buffer *pending; - LL_FOREACH(manager_conn->transfer_queue, pending) { - if (ObjectID_equal(pending->object_id, obj_id) && - (pending->type == MessageType_PlasmaDataReply)) { - return; - } + PlasmaRequestBuffer *pending; + HASH_FIND(hh, manager_conn->pending_object_transfers, &obj_id, sizeof(obj_id), + pending); + if (pending != NULL) { + return; } /* If we already have a connection to this manager and its inactive, @@ -771,8 +792,8 @@ void process_transfer_request(event_loop *loop, counter += 1; } while (obj_buffer.data_size == -1); DCHECK(obj_buffer.metadata == obj_buffer.data + obj_buffer.data_size); - plasma_request_buffer *buf = - (plasma_request_buffer *) malloc(sizeof(plasma_request_buffer)); + PlasmaRequestBuffer *buf = + (PlasmaRequestBuffer *) malloc(sizeof(PlasmaRequestBuffer)); buf->type = MessageType_PlasmaDataReply; buf->object_id = obj_id; /* We treat buf->data as a pointer to the concatenated data and metadata, so @@ -781,7 +802,9 @@ void process_transfer_request(event_loop *loop, buf->data_size = obj_buffer.data_size; buf->metadata_size = obj_buffer.metadata_size; - LL_APPEND(manager_conn->transfer_queue, buf); + DL_APPEND(manager_conn->transfer_queue, buf); + HASH_ADD(hh, manager_conn->pending_object_transfers, object_id, + sizeof(buf->object_id), buf); } /** @@ -803,8 +826,8 @@ void process_data_request(event_loop *loop, int64_t data_size, int64_t metadata_size, ClientConnection *conn) { - plasma_request_buffer *buf = - (plasma_request_buffer *) malloc(sizeof(plasma_request_buffer)); + PlasmaRequestBuffer *buf = + (PlasmaRequestBuffer *) malloc(sizeof(PlasmaRequestBuffer)); buf->object_id = object_id; buf->data_size = data_size; buf->metadata_size = metadata_size; @@ -819,7 +842,7 @@ void process_data_request(event_loop *loop, if (error_code == PlasmaError_OK) { /* Add buffer where the fetched data is to be stored to * conn->transfer_queue. */ - LL_APPEND(conn->transfer_queue, buf); + DL_APPEND(conn->transfer_queue, buf); } CHECK(conn->cursor == 0); @@ -841,9 +864,9 @@ void process_data_request(event_loop *loop, } } -void request_transfer_from(plasma_manager_state *manager_state, +void request_transfer_from(PlasmaManagerState *manager_state, ObjectID object_id) { - fetch_request *fetch_req; + FetchRequest *fetch_req; HASH_FIND(hh, manager_state->fetch_requests, &object_id, sizeof(object_id), fetch_req); /* TODO(rkn): This probably can be NULL so we should remove this check, and @@ -871,8 +894,8 @@ void request_transfer_from(plasma_manager_state *manager_state, LOG_FATAL("This manager is attempting to request a transfer from itself."); } - plasma_request_buffer *transfer_request = - (plasma_request_buffer *) malloc(sizeof(plasma_request_buffer)); + PlasmaRequestBuffer *transfer_request = + (PlasmaRequestBuffer *) malloc(sizeof(PlasmaRequestBuffer)); transfer_request->type = MessageType_PlasmaDataRequest; transfer_request->object_id = fetch_req->object_id; @@ -883,16 +906,16 @@ void request_transfer_from(plasma_manager_state *manager_state, send_queued_request, manager_conn); } /* Add this transfer request to this connection's transfer queue. */ - LL_APPEND(manager_conn->transfer_queue, transfer_request); + DL_APPEND(manager_conn->transfer_queue, transfer_request); /* On the next attempt, try the next manager in manager_vector. */ fetch_req->next_manager += 1; fetch_req->next_manager %= fetch_req->manager_count; } int fetch_timeout_handler(event_loop *loop, timer_id id, void *context) { - plasma_manager_state *manager_state = (plasma_manager_state *) context; + PlasmaManagerState *manager_state = (PlasmaManagerState *) context; /* Loop over the fetch requests and reissue the requests. */ - fetch_request *fetch_req, *tmp; + FetchRequest *fetch_req, *tmp; HASH_ITER(hh, manager_state->fetch_requests, fetch_req, tmp) { if (fetch_req->manager_count > 0) { request_transfer_from(manager_state, fetch_req->object_id); @@ -901,8 +924,8 @@ int fetch_timeout_handler(event_loop *loop, timer_id id, void *context) { return MANAGER_TIMEOUT; } -bool is_object_local(plasma_manager_state *state, ObjectID object_id) { - available_object *entry; +bool is_object_local(PlasmaManagerState *state, ObjectID object_id) { + AvailableObject *entry; HASH_FIND(hh, state->local_available_objects, &object_id, sizeof(object_id), entry); return entry != NULL; @@ -912,11 +935,11 @@ void request_transfer(ObjectID object_id, int manager_count, const char *manager_vector[], void *context) { - plasma_manager_state *manager_state = (plasma_manager_state *) context; + PlasmaManagerState *manager_state = (PlasmaManagerState *) context; /* This callback is called from object_table_subscribe, which guarantees that * the manager vector contains at least one element. */ CHECK(manager_count >= 1); - fetch_request *fetch_req; + FetchRequest *fetch_req; HASH_FIND(hh, manager_state->fetch_requests, &object_id, sizeof(object_id), fetch_req); @@ -962,8 +985,8 @@ void call_request_transfer(ObjectID object_id, int manager_count, const char *manager_vector[], void *context) { - plasma_manager_state *manager_state = (plasma_manager_state *) context; - fetch_request *fetch_req; + PlasmaManagerState *manager_state = (PlasmaManagerState *) context; + FetchRequest *fetch_req; /* Check that there isn't already a fetch request for this object. */ HASH_FIND(hh, manager_state->fetch_requests, &object_id, sizeof(object_id), fetch_req); @@ -983,7 +1006,7 @@ void object_present_callback(ObjectID object_id, int manager_count, const char *manager_vector[], void *context) { - plasma_manager_state *manager_state = (plasma_manager_state *) context; + PlasmaManagerState *manager_state = (PlasmaManagerState *) context; /* This callback is called from object_table_subscribe, which guarantees that * the manager vector contains at least one element. */ CHECK(manager_count >= 1); @@ -1000,9 +1023,9 @@ void object_table_subscribe_callback(ObjectID object_id, int manager_count, const char *manager_vector[], void *context) { - plasma_manager_state *manager_state = (plasma_manager_state *) context; + PlasmaManagerState *manager_state = (PlasmaManagerState *) context; /* Run the callback for fetch requests if there is a fetch request. */ - fetch_request *fetch_req; + FetchRequest *fetch_req; HASH_FIND(hh, manager_state->fetch_requests, &object_id, sizeof(object_id), fetch_req); if (fetch_req != NULL) { @@ -1015,7 +1038,7 @@ void object_table_subscribe_callback(ObjectID object_id, void process_fetch_requests(ClientConnection *client_conn, int num_object_ids, ObjectID object_ids[]) { - plasma_manager_state *manager_state = client_conn->manager_state; + PlasmaManagerState *manager_state = client_conn->manager_state; int num_object_ids_to_request = 0; /* This is allocating more space than necessary, but we do not know the exact @@ -1032,7 +1055,7 @@ void process_fetch_requests(ClientConnection *client_conn, } /* Check if this object is already being fetched. If so, do nothing. */ - fetch_request *entry; + FetchRequest *entry; HASH_FIND(hh, manager_state->fetch_requests, &obj_id, sizeof(obj_id), entry); if (entry != NULL) { @@ -1062,7 +1085,7 @@ void process_fetch_requests(ClientConnection *client_conn, } int wait_timeout_handler(event_loop *loop, timer_id id, void *context) { - wait_request *wait_req = (wait_request *) context; + WaitRequest *wait_req = (WaitRequest *) context; return_from_wait(wait_req->client_conn->manager_state, wait_req); return EVENT_LOOP_TIMER_DONE; } @@ -1073,11 +1096,11 @@ void process_wait_request(ClientConnection *client_conn, uint64_t timeout_ms, int num_ready_objects) { CHECK(client_conn != NULL); - plasma_manager_state *manager_state = client_conn->manager_state; + PlasmaManagerState *manager_state = client_conn->manager_state; /* Create a wait request for this object. */ - wait_request *wait_req = (wait_request *) malloc(sizeof(wait_request)); - memset(wait_req, 0, sizeof(wait_request)); + WaitRequest *wait_req = (WaitRequest *) malloc(sizeof(WaitRequest)); + memset(wait_req, 0, sizeof(WaitRequest)); wait_req->client_conn = client_conn; wait_req->timer = -1; wait_req->num_object_requests = num_object_requests; @@ -1226,10 +1249,10 @@ void process_status_request(ClientConnection *client_conn, ObjectID object_id) { request_status_done, client_conn); } -void process_delete_object_notification(plasma_manager_state *state, +void process_delete_object_notification(PlasmaManagerState *state, ObjectInfo object_info) { ObjectID obj_id = object_info.obj_id; - available_object *entry; + AvailableObject *entry; HASH_FIND(hh, state->local_available_objects, &obj_id, sizeof(obj_id), entry); if (entry != NULL) { HASH_DELETE(hh, state->local_available_objects, entry); @@ -1247,11 +1270,10 @@ void process_delete_object_notification(plasma_manager_state *state, * up-to-date. */ } -void process_add_object_notification(plasma_manager_state *state, +void process_add_object_notification(PlasmaManagerState *state, ObjectInfo object_info) { ObjectID obj_id = object_info.obj_id; - available_object *entry = - (available_object *) malloc(sizeof(available_object)); + AvailableObject *entry = (AvailableObject *) malloc(sizeof(AvailableObject)); entry->object_id = obj_id; HASH_ADD(hh, state->local_available_objects, object_id, sizeof(ObjectID), entry); @@ -1266,7 +1288,7 @@ void process_add_object_notification(plasma_manager_state *state, } /* If we were trying to fetch this object, finish up the fetch request. */ - fetch_request *fetch_req; + FetchRequest *fetch_req; HASH_FIND(hh, state->fetch_requests, &obj_id, sizeof(obj_id), fetch_req); if (fetch_req != NULL) { remove_fetch_request(state, fetch_req); @@ -1284,7 +1306,7 @@ void process_object_notification(event_loop *loop, int client_sock, void *context, int events) { - plasma_manager_state *state = (plasma_manager_state *) context; + PlasmaManagerState *state = (PlasmaManagerState *) context; ObjectInfo object_info; /* Read the notification from Plasma. */ int error = @@ -1395,9 +1417,11 @@ ClientConnection *ClientConnection_init(event_loop *loop, /* Create a new data connection context per client. */ ClientConnection *conn = (ClientConnection *) malloc(sizeof(ClientConnection)); - conn->manager_state = (plasma_manager_state *) context; + conn->manager_state = (PlasmaManagerState *) context; conn->cursor = 0; conn->transfer_queue = NULL; + /* TODO(rkn): Is this pending_object_transfers hash table ever used? */ + conn->pending_object_transfers = NULL; conn->fd = new_socket; conn->active_objects = NULL; conn->num_return_objects = 0; @@ -1438,8 +1462,8 @@ void start_server(const char *store_socket_name, CHECKM(local_sock >= 0, "Unable to bind local manager socket"); g_manager_state = - init_plasma_manager_state(store_socket_name, manager_socket_name, - master_addr, port, db_addr, db_port); + PlasmaManagerState_init(store_socket_name, manager_socket_name, + master_addr, port, db_addr, db_port); CHECK(g_manager_state); CHECK(listen(remote_sock, 5) != -1); diff --git a/src/plasma/plasma_manager.h b/src/plasma/plasma_manager.h index cee07585a..ccf68bc17 100644 --- a/src/plasma/plasma_manager.h +++ b/src/plasma/plasma_manager.h @@ -20,15 +20,15 @@ /* The buffer size in bytes. Data will get transfered in multiples of this */ #define BUFSIZE 4096 -typedef struct plasma_manager_state plasma_manager_state; +typedef struct PlasmaManagerState PlasmaManagerState; typedef struct ClientConnection ClientConnection; -typedef struct client_object_request client_object_request; +typedef struct ClientObjectRequest ClientObjectRequest; /** * Initializes the plasma manager state. This connects the manager to the local * plasma store, starts the manager listening for client connections, and * connects the manager to a database if there is one. The returned manager - * state should be freed using the provided destroy_plasma_manager_state + * state should be freed using the provided PlasmaManagerState_destroy * function. * * @param store_socket_name The socket name used to connect to the local store. @@ -41,12 +41,12 @@ typedef struct client_object_request client_object_request; * @param db_port The IP port of the database to connect to. * @return A pointer to the initialized plasma manager state. */ -plasma_manager_state *init_plasma_manager_state(const char *store_socket_name, - const char *manager_socket_name, - const char *manager_addr, - int manager_port, - const char *db_addr, - int db_port); +PlasmaManagerState *PlasmaManagerState_init(const char *store_socket_name, + const char *manager_socket_name, + const char *manager_addr, + int manager_port, + const char *db_addr, + int db_port); /** * Destroys the plasma manager state and its connections. @@ -54,7 +54,7 @@ plasma_manager_state *init_plasma_manager_state(const char *store_socket_name, * @param state A pointer to the plasma manager state to destroy. * @return Void. */ -void destroy_plasma_manager_state(plasma_manager_state *state); +void PlasmaManagerState_free(PlasmaManagerState *state); /** * Process a request from another object store manager to transfer an object. @@ -167,8 +167,7 @@ ClientConnection *ClientConnection_init(event_loop *loop, */ /* Buffer for requests between plasma managers. */ -typedef struct plasma_request_buffer plasma_request_buffer; -struct plasma_request_buffer { +typedef struct PlasmaRequestBuffer { int type; ObjectID object_id; uint8_t *data; @@ -178,8 +177,16 @@ struct plasma_request_buffer { /* Pointer to the next buffer that we will write to this plasma manager. This * field is only used if we're pushing requests to another plasma manager, * not if we are receiving data. */ - plasma_request_buffer *next; -}; + PlasmaRequestBuffer *next; + /* This is required to implement a doubly-linked list. We do not use this + * field except through the UT_list macros. */ + PlasmaRequestBuffer *prev; + /* This is used to also store the PlasmaRequestBuffer in a hash table. The + * hash table is used as a set to make sure we don't try to send the same + * object multiple times to the same manager. The object_id field will be the + * key to the hash table. */ + UT_hash_handle hh; +} PlasmaRequestBuffer; /** * Call the request_transfer method, which well attempt to get an object from @@ -215,7 +222,7 @@ int fetch_timeout_handler(event_loop *loop, timer_id id, void *context); * @return Void. */ void remove_object_request(ClientConnection *client_conn, - client_object_request *object_req); + ClientObjectRequest *object_req); /** * Get a connection to the remote manager at the specified address. Creates a @@ -226,7 +233,7 @@ void remove_object_request(ClientConnection *client_conn, * @param port The port that the remote manager is listening on. * @return A pointer to the connection to the remote manager. */ -ClientConnection *get_manager_connection(plasma_manager_state *state, +ClientConnection *get_manager_connection(PlasmaManagerState *state, const char *ip_addr, int port); @@ -240,7 +247,7 @@ ClientConnection *get_manager_connection(plasma_manager_state *state, * object. 1 means that the client has sent all the data, 0 means there * is more. */ -int read_object_chunk(ClientConnection *conn, plasma_request_buffer *buf); +int read_object_chunk(ClientConnection *conn, PlasmaRequestBuffer *buf); /** * Writes an object chunk from a buffer to the given client. This is the @@ -250,7 +257,7 @@ int read_object_chunk(ClientConnection *conn, plasma_request_buffer *buf); * @param buf The buffer to read data from. * @return Void. */ -void write_object_chunk(ClientConnection *conn, plasma_request_buffer *buf); +void write_object_chunk(ClientConnection *conn, PlasmaRequestBuffer *buf); /** * Get the event loop of the given plasma manager state. @@ -258,7 +265,7 @@ void write_object_chunk(ClientConnection *conn, plasma_request_buffer *buf); * @param state The state of the plasma manager whose loop we want. * @return A pointer to the manager's event loop. */ -event_loop *get_event_loop(plasma_manager_state *state); +event_loop *get_event_loop(PlasmaManagerState *state); /** * Get the file descriptor for the given client's socket. This is the socket @@ -277,6 +284,6 @@ int get_client_sock(ClientConnection *conn); * @return A bool that is true if the requested object is local and false * otherwise. */ -bool is_object_local(plasma_manager_state *state, ObjectID object_id); +bool is_object_local(PlasmaManagerState *state, ObjectID object_id); #endif /* PLASMA_MANAGER_H */ diff --git a/src/plasma/test/manager_tests.cc b/src/plasma/test/manager_tests.cc index 2bee92920..642171d69 100644 --- a/src/plasma/test/manager_tests.cc +++ b/src/plasma/test/manager_tests.cc @@ -46,7 +46,7 @@ typedef struct { int manager_local_fd; int local_store; int manager; - plasma_manager_state *state; + PlasmaManagerState *state; event_loop *loop; /* Accept a connection from the local manager on the remote manager. */ ClientConnection *write_conn; @@ -67,9 +67,9 @@ plasma_mock *init_plasma_mock(plasma_mock *remote_mock) { CHECK(mock->manager_local_fd >= 0 && mock->local_store >= 0); - mock->state = init_plasma_manager_state(plasma_store_socket_name, - utstring_body(manager_socket_name), - manager_addr, mock->port, NULL, 0); + mock->state = PlasmaManagerState_init(plasma_store_socket_name, + utstring_body(manager_socket_name), + manager_addr, mock->port, NULL, 0); mock->loop = get_event_loop(mock->state); /* Accept a connection from the local manager on the remote manager. */ if (remote_mock != NULL) { @@ -99,7 +99,7 @@ void destroy_plasma_mock(plasma_mock *mock) { close(get_client_sock(mock->read_conn)); free(mock->read_conn); } - destroy_plasma_manager_state(mock->state); + PlasmaManagerState_free(mock->state); free(mock->client_conn); plasma_disconnect(mock->plasma_conn); close(mock->local_store); @@ -217,14 +217,14 @@ TEST read_write_object_chunk_test(void) { const char *data = "Hello world!"; const int data_size = strlen(data) + 1; const int metadata_size = 0; - plasma_request_buffer remote_buf; + PlasmaRequestBuffer remote_buf; remote_buf.type = MessageType_PlasmaDataReply; remote_buf.object_id = oid; remote_buf.data = (uint8_t *) data; remote_buf.data_size = data_size; remote_buf.metadata = (uint8_t *) data + data_size; remote_buf.metadata_size = metadata_size; - plasma_request_buffer local_buf; + PlasmaRequestBuffer local_buf; local_buf.object_id = oid; local_buf.data_size = data_size; local_buf.metadata_size = metadata_size;