Triton-commits
Threads by month
- ----- 2026 -----
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2025 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2024 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2023 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2022 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2021 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2020 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2019 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2018 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2017 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2016 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2015 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2014 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2013 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2012 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2011 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2010 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2009 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
August 2012
- 1 participants
- 48 discussions
31 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via 5c990ab54e21e699f6925ed9469d8bc009f05327 (commit)
from 340ce550471845578fb757c6e55931cc0f1bd486 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit 5c990ab54e21e699f6925ed9469d8bc009f05327
Author: Judicael Zounmevo <zounmevoj(a)mcs.anl.gov>
Date: Fri Aug 31 14:26:59 2012 -0500
Update to the ssm net-module and minors improvements
-Added match entries cleanup before an SSM instance is stopped
-Delayed shutdown phase if network operations are still pending
-Forbid further network operation initiation when shutdown is in progress
-Added environment variables to specify an initial ethernet port value
-Made some other fixes
-Added support for an ip address file in the protocol test over SSM
-Made the internal MPI one-sided buffer pool size parametrable
through environment variable
-----------------------------------------------------------------------
Summary of changes:
code/src/common/triton-error.spec | 1 +
code/src/net/mpi/mpi-method.ae | 163 +++++++++++++++--------------
code/src/net/ssm/ssm-method.ae | 198 +++++++++++++++++++++++++----------
code/src/net/tests/ssm-one-sided.ae | 88 +++++++++++++---
code/src/net/tests/ssm-sr-test.ae | 25 ++++-
5 files changed, 322 insertions(+), 153 deletions(-)
Diff of changes:
diff --git a/code/src/common/triton-error.spec b/code/src/common/triton-error.spec
index c3c94ba..972dc54 100644
--- a/code/src/common/triton-error.spec
+++ b/code/src/common/triton-error.spec
@@ -66,3 +66,4 @@ ERR_MEM_REGISTERED "The memory is already registered"
ERR_SSM_UNKNOWN "Unknown SSM error"
ERR_SSM_TRSPT_TYPE_UNAVAILABLE "SSM transport type unavailable at this endpoint"
ERR_NAME_HASH_COLLISION "Name hash collision"
+ERR_NET_MODULE_SHUTTING_DOWN "Operation forbidden because net layer is shutting down"
diff --git a/code/src/net/mpi/mpi-method.ae b/code/src/net/mpi/mpi-method.ae
index c2a32c4..5644c3f 100644
--- a/code/src/net/mpi/mpi-method.ae
+++ b/code/src/net/mpi/mpi-method.ae
@@ -11,7 +11,7 @@
#include "src/common/triton-init.h"
#include "src/remote/byteswap.h"
-/* -----------------------------<Some one-sided stuff>---------------------------------------------*/
+
#define PRINT_ERR_MSG(err_msg_format,...) \
do \
{ \
@@ -43,23 +43,20 @@ do
#define KB 1024
#define MB 1048576
-#define TRITON_MPI_MAX_USABLE_REGISTERED_MEM_SIZE (512*MB)
-#define TRITON_MPI_REGISTERED_RMA_MEM_SIZE TRITON_MPI_MAX_USABLE_REGISTERED_MEM_SIZE
+#define TRITON_DEFAULT_MPI_MAX_USABLE_REGISTERED_MEM_SIZE (512*MB)
#define ZERO_SIZE_MEM_HANDLE 0x80000000CAFEBABE
#define WORD_SIZE SIZEOF_VOID_P
-#if WORD_SIZE==4
-#elif WORD_SIZE==8
-#else
-#error "Unsupported word size"
-#endif
-#define DEFAULT_MPI_ONE_SIDED_SIZE_THRESHOLD 0
+#define DEFAULT_MPI_ONE_SIDED_SIZE_THRESHOLD 0 /*this is temporarily set to 0 for tests purposes
+ should be changed to some sane value (4K?)
+ */
static uint64_t mpi_one_sided_size_threshold = DEFAULT_MPI_ONE_SIDED_SIZE_THRESHOLD;
+static uint32_t mpi_internal_one_sided_buffer_size = TRITON_DEFAULT_MPI_MAX_USABLE_REGISTERED_MEM_SIZE;
-/*--------------------------------</Some one-sided stuff>-----------------------------------------*/
+#define ENV_INTERNAL_ONE_SIDED_BUFFER_SIZE "INTERNAL_ONE_SIDED_BUF_SIZE"
static triton_msg_method_addr_t mpi_rank;
@@ -490,8 +487,8 @@ typedef struct locked_rank_t
{
struct locked_rank_t* next;
int rank;
- int lock_count;//how many pending lock request for this remote rank. Th ecount include the currently holding
- triton_mutex_t mutex; //the same rank cannot be locked more than once by the same process
+ int lock_count;//how many pending lock request for this remote rank. The count includes the currently holding
+ triton_mutex_t mutex; //because the same rank cannot be locked more than once by the same process
}locked_rank_t;
struct
@@ -507,12 +504,12 @@ struct
locked_rank_t* locked_rank_head;
locked_rank_t* locked_rank_tail;
-}global_data = {-1, -1, NULL, NULL, 0, TRITON_MUTEX_INITIALIZER, NULL, NULL, NULL};
+}mpi_global_data = {-1, -1, NULL, NULL, 0, TRITON_MUTEX_INITIALIZER, NULL, NULL, NULL};
static locked_rank_t* find_locked_rank(int rank)
{
- locked_rank_t* lr = global_data.locked_rank_head;
+ locked_rank_t* lr = mpi_global_data.locked_rank_head;
while(lr)
{
if(lr->rank == rank)
@@ -532,7 +529,7 @@ static inline void remove_locked_rank(int rank)
locked_rank_t* lr = find_locked_rank(rank);
if(lr)
{
- remove_node((void**)(&global_data.locked_rank_head), (void**)(&global_data.locked_rank_tail),(void*)lr);
+ remove_node((void**)(&mpi_global_data.locked_rank_head), (void**)(&mpi_global_data.locked_rank_tail),(void*)lr);
triton_mutex_destroy(&lr->mutex);
free(lr);
}
@@ -545,7 +542,7 @@ static locked_rank_t* add_locked_rank(int rank)
lr->rank = rank;
lr->lock_count = 0;
triton_mutex_init(&lr->mutex, NULL);
- add_node((void**)(&global_data.locked_rank_head), (void**)(&global_data.locked_rank_tail),(void*)lr);
+ add_node((void**)(&mpi_global_data.locked_rank_head), (void**)(&mpi_global_data.locked_rank_tail),(void*)lr);
return lr;
}
@@ -555,7 +552,7 @@ static locked_rank_t* add_locked_rank(int rank)
static mem_range_t* find_associated_mem_range(void* address)
{
mem_range_t* cur_mem_range_handle = NULL;
- for(cur_mem_range_handle = global_data.mem_range_head;
+ for(cur_mem_range_handle = mpi_global_data.mem_range_head;
cur_mem_range_handle; cur_mem_range_handle = cur_mem_range_handle->next)
{
if(cur_mem_range_handle->mem_reg_info &&
@@ -592,15 +589,15 @@ static triton_ret_t register_memory_internal(void* base,
}
//Check for space availability
- if(size > global_data.remaining_mem_size)
+ if(size > mpi_global_data.remaining_mem_size)
{
return TRITON_ERR_NOMEM;
}
- triton_mutex_lock(&global_data.mutex);
+ triton_mutex_lock(&mpi_global_data.mutex);
//check if any contiguous memory slice exists that is at least size
- cur_mem_range_handle = global_data.mem_range_head;
+ cur_mem_range_handle = mpi_global_data.mem_range_head;
while(cur_mem_range_handle)
{
if(cur_mem_range_handle->mem_reg_info == NULL) //the range is not currently registered
@@ -612,13 +609,13 @@ static triton_ret_t register_memory_internal(void* base,
mem_range_t* new_mem_range_handle = (mem_range_t*)malloc(sizeof(mem_range_t));
if(new_mem_range_handle == NULL)
{
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_ERR_NOMEM;
}
mem_reg_info = (mem_reg_t*)malloc(sizeof(mem_reg_t));
if(mem_reg_info == NULL)
{
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_ERR_NOMEM;
}
mem_reg_info->publish_count = 0;
@@ -632,7 +629,7 @@ static triton_ret_t register_memory_internal(void* base,
if(prev_mem_range_handle)
prev_mem_range_handle->next = new_mem_range_handle;
else
- global_data.mem_range_head = new_mem_range_handle;
+ mpi_global_data.mem_range_head = new_mem_range_handle;
new_mem_range_handle->next = cur_mem_range_handle;
found_mem_range = new_mem_range_handle;
@@ -644,7 +641,7 @@ static triton_ret_t register_memory_internal(void* base,
mem_reg_info = malloc(sizeof(mem_reg_t));
if(mem_reg_info== NULL)
{
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_ERR_NOMEM;
}
mem_reg_info->publish_count = 0;
@@ -657,26 +654,26 @@ static triton_ret_t register_memory_internal(void* base,
if(found_mem_range != NULL)
{
- global_data.remaining_mem_size -= size;
+ mpi_global_data.remaining_mem_size -= size;
found_mem_range->mem_reg_info = mem_reg_info;
- mem_reg_info->base_address = (void*)((uint64_t)global_data.one_sided_memory +
+ mem_reg_info->base_address = (void*)((uint64_t)mpi_global_data.one_sided_memory +
found_mem_range->offset);
mem_reg_info->user_address = base; /*base is NULL if this function is called from buffer_allocate.
*/
mem_reg_info->user_size = user_size;
mem_reg_info->size = size;
*handle = mem_reg_info;
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_SUCCESS;
}
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_ERR_NOMEM;
}
//Just as the routine name says ...
static void compact_free_contiguous_mem_ranges()
{
- mem_range_t* cur_mem_range_handle = global_data.mem_range_head;
+ mem_range_t* cur_mem_range_handle = mpi_global_data.mem_range_head;
mem_range_t* prev_mem_range_handle = NULL;
mem_range_t* temp_mem_range_handle = NULL;
@@ -699,7 +696,7 @@ static void compact_free_contiguous_mem_ranges()
static triton_ret_t unregister_memory_internal(mem_reg_t* mem_reg_info)
{
- mem_range_t* cur_mem_range_handle = global_data.mem_range_head;
+ mem_range_t* cur_mem_range_handle = mpi_global_data.mem_range_head;
if(!mem_reg_info)
return TRITON_ERR_INVAL;
@@ -707,7 +704,7 @@ static triton_ret_t unregister_memory_internal(mem_reg_t* mem_reg_info)
if(mem_reg_info->publish_count != 0)
return TRITON_ERR_MEM_IN_USE;
- triton_mutex_lock(&global_data.mutex);
+ triton_mutex_lock(&mpi_global_data.mutex);
if(mem_reg_info->base_address == mem_reg_info->user_address)
{
/*The handle is asociated with an address created through buffer_allocate.
@@ -716,23 +713,23 @@ static triton_ret_t unregister_memory_internal(mem_reg_t* mem_reg_info)
(See the comment on mem_reg_t to understand the comment above)
*/
mem_reg_info->user_address = NULL;
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_SUCCESS;
}
while(cur_mem_range_handle)
{
if(cur_mem_range_handle->mem_reg_info == mem_reg_info)
{
- global_data.remaining_mem_size += mem_reg_info->size;
+ mpi_global_data.remaining_mem_size += mem_reg_info->size;
free(cur_mem_range_handle->mem_reg_info);
cur_mem_range_handle->mem_reg_info = NULL;
compact_free_contiguous_mem_ranges();
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_SUCCESS;
}
cur_mem_range_handle = cur_mem_range_handle->next;
}
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_ERR_INVAL; //The handle is not a valid registered memory handle.
}
@@ -807,7 +804,7 @@ static triton_ret_t unregister_memory(triton_mem_reg_handle_t handle /*handle pr
static int is_valid_mem_reg_info(mem_reg_t* mem_reg_info)
{
- mem_range_t* cur_mem_range_handle = global_data.mem_range_head;
+ mem_range_t* cur_mem_range_handle = mpi_global_data.mem_range_head;
while(cur_mem_range_handle)
{
@@ -850,14 +847,14 @@ static __blocking triton_ret_t publish_memory(triton_mem_reg_handle_t in_handle,
return TRITON_ERR_NOMEM;
mem_pub_info->offset = offset + ((uint64_t)mem_reg_info->base_address -
- (uint64_t)global_data.one_sided_memory);
+ (uint64_t)mpi_global_data.one_sided_memory);
mem_pub_info->size = size;
mem_pub_info->mem_reg_info = mem_reg_info;
- mem_pub_info->client_rank = global_data.rank;
+ mem_pub_info->client_rank = mpi_global_data.rank;
- triton_mutex_lock(&global_data.mutex);
+ triton_mutex_lock(&mpi_global_data.mutex);
mem_reg_info->publish_count++;
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
//if the user-address has not been allocated using buffer_allocate
if(mem_reg_info->base_address != mem_reg_info->user_address)
@@ -866,7 +863,7 @@ static __blocking triton_ret_t publish_memory(triton_mem_reg_handle_t in_handle,
that any one-sided seamingly reading from the user memory could
effectively grab data put by the user.
*/
- memcpy((void*)((uint64_t)global_data.one_sided_memory + mem_pub_info->offset),
+ memcpy((void*)((uint64_t)mpi_global_data.one_sided_memory + mem_pub_info->offset),
(void*)((uint64_t)mem_reg_info->user_address + mem_pub_info->offset),
mem_pub_info->size);
}
@@ -897,15 +894,15 @@ static __blocking triton_ret_t unpublish_memory(triton_mem_pub_handle_t handle /
if(mem_reg_info->base_address != mem_reg_info->user_address)
{
memcpy((void*)((uint64_t)mem_reg_info->user_address + mem_pub_info->offset),
- (void*)((uint64_t)global_data.one_sided_memory + mem_pub_info->offset),
+ (void*)((uint64_t)mpi_global_data.one_sided_memory + mem_pub_info->offset),
mem_pub_info->size);
}
free(mem_pub_info);
- triton_mutex_lock(&global_data.mutex);
+ triton_mutex_lock(&mpi_global_data.mutex);
mem_reg_info->publish_count--;
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_SUCCESS;
}
@@ -1012,7 +1009,7 @@ static __blocking triton_ret_t get(triton_msg_method_addr_t from,
if(iov_items_count == 1)
{
ret = MPI_Get(dest_bufs[0], iovs[0].size, MPI_CHAR, (int)from, iovs[0].offset,
- iovs[0].size, MPI_CHAR, global_data.wins[(int)from]);
+ iovs[0].size, MPI_CHAR, mpi_global_data.wins[(int)from]);
RETURN_ON_MPI_FAILURE(ret, "MPI_Get failed");
return TRITON_SUCCESS;
}
@@ -1063,7 +1060,7 @@ static __blocking triton_ret_t get(triton_msg_method_addr_t from,
}
ret = MPI_Get(MPI_BOTTOM, 1, origin_hind_type, (int)from, 0,
- 1, target_hind_type, global_data.wins[(int)from]);
+ 1, target_hind_type, mpi_global_data.wins[(int)from]);
MPI_Type_free(&target_hind_type);
MPI_Type_free(&origin_hind_type);
@@ -1091,7 +1088,7 @@ static __blocking triton_ret_t put(triton_msg_method_addr_t to,
if(iov_items_count == 1)
{
ret = MPI_Put(src_bufs[0], iovs[0].size, MPI_CHAR, (int)to, iovs[0].offset,
- iovs[0].size, MPI_CHAR, global_data.wins[(int)to]);
+ iovs[0].size, MPI_CHAR, mpi_global_data.wins[(int)to]);
RETURN_ON_MPI_FAILURE(ret, "MPI_Put failed");
return TRITON_SUCCESS;
}
@@ -1142,7 +1139,7 @@ static __blocking triton_ret_t put(triton_msg_method_addr_t to,
}
ret = MPI_Put(MPI_BOTTOM, 1, origin_hind_type, (int)to, 0,
- 1, target_hind_type, global_data.wins[(int)to]);
+ 1, target_hind_type, mpi_global_data.wins[(int)to]);
MPI_Type_free(&target_hind_type);
MPI_Type_free(&origin_hind_type);
@@ -1157,15 +1154,15 @@ static __blocking triton_ret_t start_epoch(triton_mem_pub_handle_t handle)
mem_pub_t *mem_pub = (mem_pub_t*)handle;
int ret;
locked_rank_t* lr = NULL;
- triton_mutex_lock(&global_data.mutex);
+ triton_mutex_lock(&mpi_global_data.mutex);
lr = find_locked_rank(mem_pub->client_rank);
if(!lr)
lr = add_locked_rank(mem_pub->client_rank);
lr->lock_count++;
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
triton_mutex_lock(&lr->mutex);
ret = MPI_Win_lock(MPI_LOCK_SHARED, mem_pub->client_rank, 0,
- global_data.wins[mem_pub->client_rank]);
+ mpi_global_data.wins[mem_pub->client_rank]);
RETURN_ON_MPI_FAILURE(ret, "MPI_Win_lock failed");
return TRITON_SUCCESS;
}
@@ -1177,17 +1174,17 @@ static __blocking triton_ret_t end_epoch(triton_mem_pub_handle_t handle)
int ret;
locked_rank_t* lr = NULL;
- ret = MPI_Win_unlock(mem_pub->client_rank, global_data.wins[mem_pub->client_rank]);
+ ret = MPI_Win_unlock(mem_pub->client_rank, mpi_global_data.wins[mem_pub->client_rank]);
RETURN_ON_MPI_FAILURE(ret, "MPI_Win_unlock failed");
- triton_mutex_lock(&global_data.mutex);
+ triton_mutex_lock(&mpi_global_data.mutex);
lr = find_locked_rank(mem_pub->client_rank);
triton_assert(lr); //the lr MUST exist here!
lr->lock_count--;
triton_mutex_unlock(&lr->mutex);
if(lr->lock_count == 0)
remove_locked_rank(mem_pub->client_rank);
- triton_mutex_unlock(&global_data.mutex);
+ triton_mutex_unlock(&mpi_global_data.mutex);
return TRITON_SUCCESS;
}
@@ -1240,33 +1237,45 @@ triton_ret_t triton_msg_mpi_init(void)
{
int ret;
int i;
+ int env_value_internal_buf_size;
+ char* env_value_str;
+
triton_ret_t tret;
- ret = MPI_Comm_rank(MPI_COMM_WORLD, &global_data.rank);
+ ret = MPI_Comm_rank(MPI_COMM_WORLD, &mpi_global_data.rank);
RETURN_ON_MPI_FAILURE(ret, "MPI_Comm_rank failed");
- mpi_rank = (triton_msg_method_addr_t)global_data.rank;
- ret = MPI_Comm_size(MPI_COMM_WORLD, &global_data.size);
+ mpi_rank = (triton_msg_method_addr_t)mpi_global_data.rank;
+ ret = MPI_Comm_size(MPI_COMM_WORLD, &mpi_global_data.size);
RETURN_ON_MPI_FAILURE(ret, "MPI_Comm_size failed");
+ env_value_str = getenv(ENV_INTERNAL_ONE_SIDED_BUFFER_SIZE);
+ if(env_value_str)
+ {
+ env_value_internal_buf_size = atoi(env_value_str);
+ if(env_value_internal_buf_size)
+ mpi_internal_one_sided_buffer_size = (uint32_t)env_value_internal_buf_size;
+ }
+
+
//Create the Win for one-sided; as well as the machinery for mem registration
- triton_mutex_init(&global_data.mutex, NULL);
- global_data.one_sided_memory = malloc(TRITON_MPI_REGISTERED_RMA_MEM_SIZE );
- ABORT_IF(global_data.one_sided_memory == NULL, "Memory allocation failed");
+ triton_mutex_init(&mpi_global_data.mutex, NULL);
+ mpi_global_data.one_sided_memory = malloc(mpi_internal_one_sided_buffer_size);
+ ABORT_IF(mpi_global_data.one_sided_memory == NULL, "Memory allocation failed");
- global_data.remaining_mem_size = TRITON_MPI_MAX_USABLE_REGISTERED_MEM_SIZE;
- global_data.mem_range_head = (mem_range_t*)calloc(1, sizeof(mem_range_t));
- ABORT_IF(global_data.mem_range_head == NULL, "Memory allocation failed");
+ mpi_global_data.remaining_mem_size = mpi_internal_one_sided_buffer_size;
+ mpi_global_data.mem_range_head = (mem_range_t*)calloc(1, sizeof(mem_range_t));
+ ABORT_IF(mpi_global_data.mem_range_head == NULL, "Memory allocation failed");
- global_data.mem_range_head->offset = 0;
- global_data.mem_range_head->size = TRITON_MPI_MAX_USABLE_REGISTERED_MEM_SIZE;
- global_data.mem_range_head->mem_reg_info = NULL;
+ mpi_global_data.mem_range_head->offset = 0;
+ mpi_global_data.mem_range_head->size = mpi_internal_one_sided_buffer_size;
+ mpi_global_data.mem_range_head->mem_reg_info = NULL;
- global_data.wins = malloc(sizeof(MPI_Win)*global_data.size);
- ABORT_IF(global_data.wins == NULL, "Memory allocation failed");
+ mpi_global_data.wins = malloc(sizeof(MPI_Win)*mpi_global_data.size);
+ ABORT_IF(mpi_global_data.wins == NULL, "Memory allocation failed");
- for(i=0; i<global_data.size; i++)
+ for(i=0; i<mpi_global_data.size; i++)
{
- ret = MPI_Win_create(global_data.one_sided_memory, TRITON_MPI_REGISTERED_RMA_MEM_SIZE,
- sizeof(char), MPI_INFO_NULL, MPI_COMM_WORLD, &global_data.wins[i]);
+ ret = MPI_Win_create(mpi_global_data.one_sided_memory, mpi_internal_one_sided_buffer_size,
+ sizeof(char), MPI_INFO_NULL, MPI_COMM_WORLD, &mpi_global_data.wins[i]);
ABORT_ON_MPI_FAILURE(ret, "MPI_Win_create failed");
}
@@ -1292,12 +1301,12 @@ void triton_msg_mpi_finalize(void)
triton_hash_destroy_and_finalize(group_table, struct triton_mpi_group, link, group_entry_free);
group_table = NULL;
}
- for(i=0; i<global_data.size; i++)
- MPI_Win_free(&global_data.wins[i]);
- free(global_data.one_sided_memory);
- triton_assert(global_data.mem_range_head->next == NULL);
- free(global_data.mem_range_head);
- triton_mutex_destroy(&global_data.mutex);
+ for(i=0; i<mpi_global_data.size; i++)
+ MPI_Win_free(&mpi_global_data.wins[i]);
+ free(mpi_global_data.one_sided_memory);
+ triton_assert(mpi_global_data.mem_range_head->next == NULL);
+ free(mpi_global_data.mem_range_head);
+ triton_mutex_destroy(&mpi_global_data.mutex);
}
__attribute__((constructor)) void triton_msg_mpi_init_register(void);
diff --git a/code/src/net/ssm/ssm-method.ae b/code/src/net/ssm/ssm-method.ae
index 452b7cb..4b309d2 100644
--- a/code/src/net/ssm/ssm-method.ae
+++ b/code/src/net/ssm/ssm-method.ae
@@ -30,7 +30,6 @@
#include <semaphore.h>
-
#define JZ_DEBUG
static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once)
{
@@ -46,7 +45,6 @@ static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, c
}
#ifdef JZ_DEBUG
- void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once);
#define JZ_WAIT_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0)
#define JZ_WAIT_FOR_DEBUGGER_IF(cond) \
do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0);}while(0)
@@ -109,11 +107,8 @@ static inline __blocking triton_ret_t triton_sem_count_down_no_cancel (triton_se
/*THE MATCH BITS MANAGEMENT STRATEGY:
1- The most significant bit distinguishes between 2-sided and one-sided. If set, it is 2-sided
2- For 2-sided:
- The group name is hashed over the 16 next most significant bits
- The tags which can therefore range from 0 to 0x7fff cover the 15 remaining bits
- 3- One-sided range from 0 to 0x7ffffffe
- 4- 0x7fffffff is reserved for message progression thread shutdown. It is kind of used by
- a process to message itself when it is in the shutdown process.
+ The group name is hashed over the 32 next most significant bits
+ The tags occupy the rest.
*/
@@ -130,6 +125,7 @@ static inline __blocking triton_ret_t triton_sem_count_down_no_cancel (triton_se
#define SSM_MAX_INTERNAL_TWO_SIDED_SIZE (SSM_MAX_TWO_SIDED_SIZE+MAX_RAW_ADDRESS_SIZE)
#define SSM_ANY_SOURCE 0
+
/*
GENERIC NOTES:
1- This net-module will potentially manage several transports
@@ -176,20 +172,22 @@ typedef struct network_endpoint_t
ssm_Iaddr address_interface;
ssm_Haddr local_listener_address_handle;
ssm_me shutdown_me; //the progress engine shutdown is triggered by an ssm_unlink op.
+ ssm_mr shutdown_mr;
unsigned char raw_address[MAX_RAW_ADDRESS_SIZE];
ssm_transport_abstraction_t* transport_abstraction;
ssm_transport_type_t transport_type;
- sem_t release_sem; //used in the release phase
+ sem_t release_sem; //used in the release phase before shutting down the network endpoint
triton_list_link_t list_link;
struct triton_hash_link address_hash_link; //if the address is not encodable over 56-bit, this contains
//the actual address data. It uses the address as a key then
+ //Transports like TCP do not use this hash_list
}network_endpoint_t;
network_endpoint_t tcp_network_endpoint;
-//TODO: Add network endpoints for other transports when they are available
+//TODO [FUTURE TRANSPORTS]: Add network endpoints for other transports when they are available
network_endpoint_t *default_net_endpoint;
@@ -210,17 +208,55 @@ struct
triton_list_link_t ptp_match_entries_list;
pthread_t progress_thread;
+
+ int nb_linked_match_entries;
+ triton_mutex_t match_entries_mutex;
+ sem_t match_entries_semaphore;
BOOL exiting; //TRUE when the process is shuting down
}ssm_global_data;
+static inline ssm_me triton_ssm_link(ssm_id id, ssm_bits bits, ssm_bits mask,
+ ssm_pos pos, ssm_me anchor, ssm_cb cb, ssm_Flink flags)
+{
+ triton_mutex_lock(&ssm_global_data.match_entries_mutex);
+ ++ssm_global_data.nb_linked_match_entries;
+ triton_mutex_unlock(&ssm_global_data.match_entries_mutex);
+ return ssm_link(id, bits, mask, pos, anchor, cb, flags);
+}
+
+
+static inline void update_match_entry_count_on_unlink()
+{
+ triton_mutex_lock(&ssm_global_data.match_entries_mutex);
+ --ssm_global_data.nb_linked_match_entries;
+ if(ssm_global_data.exiting && ssm_global_data.nb_linked_match_entries == 0)
+ sem_post(&ssm_global_data.match_entries_semaphore);
+ triton_mutex_unlock(&ssm_global_data.match_entries_mutex);
+}
+
+static inline int triton_ssm_unlink(ssm_id id, ssm_me me)
+{
+ int status = ssm_unlink(id, me);
+ if(status == SSM_REMOVE_OK)
+ {
+ update_match_entry_count_on_unlink();
+ }
+
+ return status;
+}
+
+#define FORBID_IF_SHUT_DOWN_IN_PROGRESS() \
+ if(ssm_global_data.exiting) return TRITON_ERR_NET_MODULE_SHUTTING_DOWN;
+
static void shutdown_callback_routine(void* cb_args, void* ev_data)
{
if(((ssm_result)ev_data)->op == SSM_OP_UNLINK)
{
sem_post(&((network_endpoint_t*)cb_args)->release_sem);
+ update_match_entry_count_on_unlink();
}
}
@@ -261,25 +297,28 @@ static void init_ethernet()
shutdown_callback.pcb = shutdown_callback_routine;
shutdown_callback.cbdata = &tcp_network_endpoint;
- tcp_network_endpoint.shutdown_me = ssm_link(tcp_network_endpoint.ssm_instance,
- PROGRESSION_SHUTDOWN_RESERVED_MATCH_BITS, 0, SSM_POS_HEAD, NULL, &shutdown_callback, SSM_NOF);
-
+ tcp_network_endpoint.shutdown_me = triton_ssm_link(tcp_network_endpoint.ssm_instance,
+ PROGRESSION_SHUTDOWN_RESERVED_MATCH_BITS, 0, SSM_POS_HEAD, NULL, &shutdown_callback,
+ SSM_LINK_AUTO_UNLINK);
+ tcp_network_endpoint.shutdown_mr = ssm_mr_create(NULL, NULL, 0);
+
sem_init(&tcp_network_endpoint.release_sem, 0, 0);
triton_list_link_clear(&tcp_network_endpoint.list_link);
triton_list_add_back(&tcp_network_endpoint.list_link, &ssm_global_data.net_endpoint_list);
}
-/*Provide the init_ethernet equivalent for the other transports as well
- */
static inline void release_ethernet()
{
- ssm_unlink(tcp_network_endpoint.ssm_instance, tcp_network_endpoint.shutdown_me);
+ ssm_drop(tcp_network_endpoint.ssm_instance, tcp_network_endpoint.shutdown_me,
+ tcp_network_endpoint.shutdown_mr);
sem_wait(&tcp_network_endpoint.release_sem);
ssm_stop(tcp_network_endpoint.ssm_instance);
sem_destroy(&tcp_network_endpoint.release_sem);
}
+/*TODO [FUTURE TRANSPORTS]: Provide the init_ethernet and release_ethernet equivalent for
+ the other transports as well */
static inline void set_default_net_endpoint()
{
@@ -289,10 +328,12 @@ static inline void set_default_net_endpoint()
static ssm_transport_abstraction_t* get_transport_abstraction(ssm_transport_type_t type)
{
- switch(SSM_TCP)
+ switch(type)
{
case SSM_TCP:
return get_tcp_abstraction();
+ default:
+ ;
}
return NULL;
}
@@ -305,10 +346,14 @@ static inline ssm_transport_abstraction_t* get_transport_abstraction_from_addr(t
static network_endpoint_t* get_network_endpoint(ssm_transport_type_t type)
{
- switch(SSM_TCP)
+ switch(type)
{
case SSM_TCP:
return &tcp_network_endpoint;
+
+ //TODO [FUTURE TRANSPORTS] ...
+ default:
+ ;
}
return NULL;
}
@@ -316,6 +361,9 @@ static network_endpoint_t* get_network_endpoint(ssm_transport_type_t type)
/*
Add a remote address that will be used afterward simply through a triton_msg_method_addr_t
The address is not added if it already exists
+ This method is used after recv_any gets some address. The address is cached so that
+ recv could use it afterward. The address is created from raw data; and if it requires
+ internal data creation, the internal data is cached
*/
static void add_address(network_endpoint_t* net_endpoint, unsigned char* raw_host_bytes_addr_data, triton_msg_method_addr_t* addr)
{
@@ -326,17 +374,17 @@ static void add_address(network_endpoint_t* net_endpoint, unsigned char* raw_hos
raw_host_bytes_addr_data, NULL);
break;
- /*The other ones will probably have some kind of data structure that will will be created and keyed
+ /*TODO [FUTURE TRANSPORTS]: The other ones will probably have some kind of data structure that will will be created and keyed
with the triton_msg_method_addr_t
*/
default:
- ;
+ ;
}
}
static __blocking triton_msg_method_addr_t ssm_addr_self(void)
{
- //TODO: How do we know what the caller of this function want if this ssm_instance of the net-module
+ //TODO: How do we know what the caller of this function wants if this ssm_instance of the net-module
//has several transports underneath (E.g. TCP, IB, UDP, some other stuff)
return (triton_msg_method_addr_t)default_net_endpoint->transport_abstraction->get_address(
default_net_endpoint->raw_address, NULL);
@@ -364,8 +412,11 @@ static triton_ret_t ssm_addr_free(triton_msg_method_addr_t addr)
{
case SSM_TCP:
break; //nothing to do for TCP; it uses the type|address enconding
+ //TODO [FUTURE TRANSPORTS] ...
+ default:
+ ;
}
- //The other ones in the future might be reference-counted and freed only when all copies are freed
+ //TODO [FUTURE TRANSPORTS]: The other ones in the future might be reference-counted and freed only when all copies are freed
return TRITON_SUCCESS;
}
@@ -375,7 +426,10 @@ static triton_ret_t ssm_addr_copy(const triton_msg_method_addr_t orig_addr, trit
{
case SSM_TCP:
*copy = orig_addr;
- break; //Other ssm_transport that maintain additional data might need reference counting
+ break;
+ //TODO [FUTURE TRANSPORTS]: Other ssm_transport that maintain additional data might need reference counting
+ default:
+ ;
}
return TRITON_SUCCESS;
}
@@ -391,7 +445,7 @@ static triton_ret_t ssm_addr_lookup(const char *name, triton_msg_method_addr_t *
strcpy(name_1, name);
*addr = (triton_msg_method_addr_t)tcp_network_endpoint.transport_abstraction->get_address_from_str_addr(name_1, NULL);
g_looked_up_addr = addr;
- }
+ } //TODO [FUTURE TRANSPORTS] ...
else
{
triton_assert("Unsupported ssm_transport!");
@@ -422,22 +476,23 @@ static ssm_Haddr tcp_get_ssm_addr_from_triton_addr(triton_msg_method_addr_t addr
return tcp_abstraction->create_addr_handle(tcp_network_endpoint.address_interface, tcp_args);
}
+//TODO [FUTURE TRANSPORTS]: ... equivalent of function above
+
static triton_ret_t get_ssm_addr_from_triton_addr(triton_msg_method_addr_t addr, ssm_Haddr* the_ssm_addr)
{
- switch(SSM_TCP)
+ switch(TRANSPORT_TYPE_FROM_ADDR(addr))
{
case SSM_TCP:
*the_ssm_addr = tcp_get_ssm_addr_from_triton_addr(addr);
return TRITON_SUCCESS;
+ //TODO [FUTURE TRANSPORTS] ...
+ default:
+ ;
}
return TRITON_ERR_SSM_TRSPT_TYPE_UNAVAILABLE;
}
-/*
- */
-
-
#define SSM_TAG_MAX TWO_SIDED_TAG_MASK
struct triton_ssm_group
{
@@ -550,6 +605,8 @@ static __blocking triton_ret_t ssm_send(
struct triton_ssm_group *group_entry;
ssm_bits match_bits;
+ FORBID_IF_SHUT_DOWN_IN_PROGRESS();
+
if(count>1)
return TRITON_ERR_NOT_IMPLEMENTED; //for now, only contiguous stuff is supported
if(sizes[0] > SSM_MAX_TWO_SIDED_SIZE)
@@ -690,7 +747,6 @@ static inline void cancel_receive(generic_recv_comm_info_t* generic_comm_info,t
static void recv_callback(void* cb_data, void* event_data)
{
static int counter = 0;
- JZ_WAIT_FOR_DEBUGGER();
TRACE("Entering recv_callback no = %d \n", counter);
int ret;
generic_recv_comm_info_t* comm_info = (generic_recv_comm_info_t*)cb_data;
@@ -720,12 +776,15 @@ static void recv_callback(void* cb_data, void* event_data)
free(mr_info.base);
return;
}
+
*specific_receive_info->received_bytes = result->bytes - MAX_RAW_ADDRESS_SIZE;
*specific_receive_info->sender = sender;
*specific_receive_info->tag = result->bits & TWO_SIDED_TAG_MASK;
- if(specific_receive_info->size <= *specific_receive_info->received_bytes)
+ TRACE("(result->bytes, specific_receive_info->size, specific_receive_info->received_bytes) = (%u, %u, %u)\n",
+ result->bytes, specific_receive_info->size, *specific_receive_info->received_bytes);
+ if(TRUE || specific_receive_info->size >= *specific_receive_info->received_bytes)
memcpy(specific_receive_info->buffer, (char*)mr_info.base + MAX_RAW_ADDRESS_SIZE,
- mr_info.span - MAX_RAW_ADDRESS_SIZE);
+ 100 /**specific_receive_info->received_bytes*/);
triton_sem_up(&specific_receive_info->semaphore);
@@ -733,12 +792,13 @@ static void recv_callback(void* cb_data, void* event_data)
if(comm_info->specific_recv_list.count == 0)
{
- ssm_unlink(comm_info->net_endpoint->ssm_instance, comm_info->match_entry);
+ triton_ssm_unlink(comm_info->net_endpoint->ssm_instance, comm_info->match_entry);
triton_list_del(&comm_info->list_link);
free(comm_info);
}
triton_mutex_unlock(&ssm_global_data.two_sided_mutex);
- TRACE("Exiting recv_callback no = %d (status, op, ssm_bits) = ( %d, %d, 0x%llx)\n", counter++, result->status, result->op, (uint64_t)result->bits);
+ TRACE("Exiting recv_callback no = %d (status, op, ssm_bits) = ( %d, %d, 0x%llx)\n",
+ counter++, result->status, result->op, (uint64_t)result->bits);
}
@@ -764,6 +824,8 @@ static __blocking triton_ret_t generic_receive(triton_msg_method_group_t group,
network_endpoint_t* net_endpoint = NULL;
char *buffer = NULL;
+ FORBID_IF_SHUT_DOWN_IN_PROGRESS();
+
specific_recv_comm_info.sender = sender;
specific_recv_comm_info.buffer = user_buffer;
specific_recv_comm_info.size = buf_size;
@@ -811,7 +873,7 @@ static __blocking triton_ret_t generic_receive(triton_msg_method_group_t group,
triton_list_add_back(&comm_info->list_link, &ssm_global_data.generic_recv_data_list);
callback.cbdata = comm_info;
comm_info->is_receive_any = is_receive_any;
- comm_info->match_entry = ssm_link(net_endpoint->ssm_instance, match_bits,
+ comm_info->match_entry = triton_ssm_link(net_endpoint->ssm_instance, match_bits,
is_receive_any?TWO_SIDED_ANY_SOURCE_MATCH_MASK:0, SSM_POS_TAIL, NULL, &callback, SSM_NOF);
g_recv_comm = comm_info;
}
@@ -842,7 +904,7 @@ static __blocking triton_ret_t generic_receive(triton_msg_method_group_t group,
}
-static __blocking triton_ret_t ssm_recv(
+static inline __blocking triton_ret_t ssm_recv(
triton_msg_method_addr_t from,
triton_msg_method_group_t group,
triton_msg_tag_t tag,
@@ -857,7 +919,7 @@ static __blocking triton_ret_t ssm_recv(
}
-static __blocking triton_ret_t ssm_recv_any(
+static inline __blocking triton_ret_t ssm_recv_any(
triton_msg_method_group_t group,
triton_msg_method_addr_t *from,
triton_msg_tag_t *tag,
@@ -944,7 +1006,7 @@ typedef struct mem_pub_t
int nb_expected_completions; //for the target
triton_mutex_t mutex;
triton_sem_t comm_semaphore; //MUST be used EXCLUSIVELY for communications. It is not a generic sem
- triton_sem_t unlink_semaphore; //MUST be used EXCLUSIVELY for unlink
+ sem_t unlink_semaphore; //MUST be used EXCLUSIVELY for unlink
ssm_me match_entry;
union
{
@@ -953,7 +1015,7 @@ typedef struct mem_pub_t
than this array
*/
ssmptcp_addrargs_t tcp_args;
- //TODO: When they become available, add other ssm_transport address_arg types in this union.
+ //TODO: [FUTURE TRANSPORTS] When they become available, add other ssm_transport address_arg types in this union.
//For now, only TCP is defined
};
int made_from_deserialization; //1 if made from deserialization
@@ -980,12 +1042,17 @@ static mem_reg_t* find_mem_reg(void* base_address)
return mem_reg;
}
+
static triton_ret_t register_memory_internal(void* base,
uint64_t size,
mem_reg_t** handle
)
{
- mem_reg_t *mem_reg = find_mem_reg(base);
+ mem_reg_t *mem_reg = NULL;
+
+ FORBID_IF_SHUT_DOWN_IN_PROGRESS();
+
+ mem_reg = find_mem_reg(base);
if(mem_reg)
{
if(mem_reg->is_registered)
@@ -1137,8 +1204,10 @@ static void target_side_wait_callback(void* callback_args, void* event_data)
if(result->op == SSM_OP_PUT || result->op == SSM_OP_GET)
triton_sem_up(&mem_pub->comm_semaphore);
else if(result->op == SSM_OP_UNLINK)
- triton_sem_up(&mem_pub->unlink_semaphore);
-
+ {
+ sem_post(&mem_pub->unlink_semaphore);
+ update_match_entry_count_on_unlink();
+ }
}
@@ -1154,6 +1223,9 @@ static __blocking triton_ret_t publish_memory(triton_mem_reg_handle_t in_handle,
)
{
int ret;
+
+ FORBID_IF_SHUT_DOWN_IN_PROGRESS();
+
if(size == 0)
{
*out_handle = ZERO_SIZE_MEM_HANDLE;
@@ -1191,19 +1263,19 @@ static __blocking triton_ret_t publish_memory(triton_mem_reg_handle_t in_handle,
mem_pub->callback.cbdata = mem_pub;
mem_pub->mr = ssm_mr_create(NULL, mem_pub->base, mem_pub->size);
- mem_pub->match_entry = ssm_link(default_net_endpoint->ssm_instance, mem_pub->match_bits, 0,
- SSM_POS_TAIL, NULL, &mem_pub->callback, SSM_NOF);
+ mem_pub->match_entry = triton_ssm_link(default_net_endpoint->ssm_instance, mem_pub->match_bits, 0,
+ SSM_POS_TAIL, NULL, &mem_pub->callback, SSM_LINK_AUTO_UNLINK);
ret = ssm_post(default_net_endpoint->ssm_instance, mem_pub->match_entry, mem_pub->mr, SSM_NOF);
if(ret)
{
- ssm_unlink(mem_pub->transport_instance, mem_pub->match_entry);
+ triton_ssm_unlink(mem_pub->transport_instance, mem_pub->match_entry);
free(mem_pub);
return TRITON_ERR_SSM_UNKNOWN;
}
triton_sem_init(&mem_pub->comm_semaphore, 0);
- triton_sem_init(&mem_pub->unlink_semaphore, 0);
+ sem_init(&mem_pub->unlink_semaphore, 0, 0);
*out_handle = (uint64_t)mem_pub;
@@ -1217,6 +1289,7 @@ static __blocking triton_ret_t unpublish_memory(triton_mem_pub_handle_t handle /
returned by a memory publishing*/
)
{
+ int status;
if(handle == ZERO_SIZE_MEM_HANDLE)
return TRITON_SUCCESS;
@@ -1226,11 +1299,14 @@ static __blocking triton_ret_t unpublish_memory(triton_mem_pub_handle_t handle /
if(mem_pub->made_from_deserialization)
return TRITON_ERR_INVAL; //this handle is not allowed to be released through unpublish
- ssm_unlink(default_net_endpoint->ssm_instance, mem_pub->match_entry);
- triton_sem_count_down_no_cancel(&mem_pub->unlink_semaphore, 1);
+ status = triton_ssm_unlink(default_net_endpoint->ssm_instance, mem_pub->match_entry);
+ if(status == SSM_REMOVE_BUSY)
+ sem_wait(&mem_pub->unlink_semaphore);
+ triton_assert(status != SSM_REMOVE_INVALID);
+
ssm_mr_destroy(mem_pub->mr);
triton_sem_destroy(&mem_pub->comm_semaphore);
- triton_sem_destroy(&mem_pub->unlink_semaphore);
+ sem_destroy(&mem_pub->unlink_semaphore);
free(mem_pub);
triton_mutex_lock(&ssm_global_data.mutex);
@@ -1353,7 +1429,7 @@ static triton_ret_t create_handle_from_serialization(triton_serialized_handle_ty
transport_abstraction->create_address_arg_data((uint8_t*)serialized_handle, mem_pub->transport_address_data);
break;
default:
- ;
+ ;
}
*((triton_mem_pub_handle_t*)pointer_to_out_handle) = (triton_mem_pub_handle_t)mem_pub;
@@ -1431,6 +1507,9 @@ static triton_ret_t issue_one_sided_op(
int i;
mem_pub_t* mem_pub = (mem_pub_t*)handle;
transaction_t* transactions = NULL;
+
+ FORBID_IF_SHUT_DOWN_IN_PROGRESS();
+
transactions = malloc(iov_items_count*sizeof(transaction_t));
if(!transactions)
return TRITON_ERR_NOMEM;
@@ -1480,6 +1559,9 @@ static __blocking triton_ret_t put(triton_msg_method_addr_t to,
static __blocking triton_ret_t start_epoch(triton_mem_pub_handle_t handle)
{
mem_pub_t *mem_pub = (mem_pub_t*)handle;
+
+ FORBID_IF_SHUT_DOWN_IN_PROGRESS();
+
network_endpoint_t *network_endpoint = get_network_endpoint(mem_pub->transport_type);
if(!network_endpoint)
return TRITON_ERR_SSM_TRSPT_TYPE_UNAVAILABLE;
@@ -1569,18 +1651,22 @@ static struct triton_msg_method triton_msg_ssm_method =
triton_ret_t triton_msg_ssm_init(void)
{
-
triton_ret_t tret;
-
triton_mutex_init(&ssm_global_data.mutex, NULL);
triton_mutex_init(&ssm_global_data.two_sided_mutex, NULL);
+ triton_mutex_init(&ssm_global_data.match_entries_mutex, NULL);
+ sem_init(&ssm_global_data.match_entries_semaphore, 0, 0);
triton_list_init(&ssm_global_data.mem_reg_list);
triton_list_init(&ssm_global_data.net_endpoint_list);
triton_list_init(&ssm_global_data.generic_recv_data_list);
+ ssm_global_data.nb_linked_match_entries = 0;
+
triton_sem_init_thread_safe(&ssm_global_data.mutex);
+ ssm_global_data.nb_linked_match_entries = 0;
+
init_ethernet();
/*Init the other available tyransports here (e.g. init_ib()*/
@@ -1605,9 +1691,6 @@ void triton_msg_ssm_finalize(void)
triton_sem_finalize_thread_safe(&ssm_global_data.mutex);
triton_msg_method_unregister("ssm");
- release_ethernet();
- /*Release the other available transports here (e.g. release_ib()*/
-
if(group_table != NULL)
{
triton_hash_destroy_and_finalize(group_table, struct triton_ssm_group, link, group_entry_free);
@@ -1615,8 +1698,15 @@ void triton_msg_ssm_finalize(void)
}
pthread_join(ssm_global_data.progress_thread, NULL);
+
+ sem_wait(&ssm_global_data.match_entries_semaphore); //make sure that no match entry is still linked
+ release_ethernet();
+ /*Release the other available transports here (e.g. release_ib()*/
+
triton_mutex_destroy(&ssm_global_data.mutex);
triton_mutex_destroy(&ssm_global_data.two_sided_mutex);
+ triton_mutex_destroy(&ssm_global_data.match_entries_mutex);
+ sem_destroy(&ssm_global_data.match_entries_semaphore);
}
__attribute__((constructor)) void triton_msg_ssm_init_register(void);
diff --git a/code/src/net/tests/ssm-one-sided.ae b/code/src/net/tests/ssm-one-sided.ae
index c073e0a..9d594d0 100644
--- a/code/src/net/tests/ssm-one-sided.ae
+++ b/code/src/net/tests/ssm-one-sided.ae
@@ -30,7 +30,6 @@ static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, c
}
#ifdef JZ_DEBUG
- void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once);
#define JZ_WAIT_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0)
#define JZ_WAIT_FOR_DEBUGGER_IF(cond) \
do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0);}while(0)
@@ -46,6 +45,7 @@ static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, c
#define ENV_ETHERNET_IP "ETHERNET_IP"
#define ENV_ETHERNET_PORT "ETHERNET_PORT"
+#define ENV_PORT_BASE "PORT_BASE"
#define ETHERNET_LOOPBACK "127.0.0.1"
#define MIN(a,b) ((a)<=(b)?(a):(b))
@@ -75,7 +75,8 @@ do
*/
-#define PORT_BASE 5000
+#define DEFAULT_PORT_BASE 5000
+uint16_t port_base = DEFAULT_PORT_BASE;
#define MAX_DATA_SIZE 50*MB
int rank = -1;
@@ -158,14 +159,62 @@ static inline int is_local_process_a_server()
return is_server(rank);
}
+
+char ip_file_name[151] = {0};
+
+
+typedef struct ip_str_t
+{
+ char ip[51];
+}ip_str_t;
+
+ip_str_t *str_ip_addresses;
+
+/*Build the IP addresses
+*/
+static triton_ret_t build_ips()
+{
+ int i;
+ str_ip_addresses = malloc(rank*sizeof(ip_str_t));
+ if(!str_ip_addresses)
+ return TRITON_ERR_NOMEM;
+ if(ip_file_name[0])
+ {
+ FILE* ip_file = fopen(ip_file_name, "r");
+ if(!ip_file)
+ {
+ PRINT("[WARNING!] Failed to open ip_file\n");
+ goto use_loopback; //fallback to using loopback.
+ //This means that if all the involved processes
+ //are not in the same node, the job will fail
+ }
+ for(i=0; i<size; i++)
+ (void)fgets(str_ip_addresses[i].ip, 51, ip_file);
+ fclose(ip_file);
+ }
+ return TRITON_SUCCESS;
+use_loopback:
+ for(i=0; i<size; i++)
+ strcpy(str_ip_addresses[i].ip, ETHERNET_LOOPBACK);
+ return TRITON_SUCCESS;
+}
+
+
static void set_environment()
{
int ret;
char env_string[101];
- sprintf(env_string, "%s=%s", ENV_ETHERNET_IP, ETHERNET_LOOPBACK);
+ char *env_value = NULL;
+ if((env_value = getenv(ENV_PORT_BASE)) != NULL)
+ {
+ int val = atoi(env_value);
+ if(val)
+ port_base = (uint16_t) val;
+ }
+ sprintf(env_string, "%s=%s", str_ip_addresses[rank].ip, str_ip_addresses[rank].ip);
ret = putenv(env_string);
triton_assert(ret == 0);
- sprintf(env_string, "%s=%d", ENV_ETHERNET_PORT, PORT_BASE+rank);
+ sprintf(env_string, "%s=%d", ENV_ETHERNET_PORT, port_base+rank);
ret = putenv(env_string);
triton_assert(ret == 0);
}
@@ -188,10 +237,8 @@ static __blocking void close_endpoints()
triton_addr_free(server_addrs[i]);
}
-/*
-Discovering servers could be implemented differently; e.g. having the
-processes read from a file who's who
-*/
+
+
static __blocking void discover_servers(void)
{
triton_ret_t tret;
@@ -203,7 +250,7 @@ static __blocking void discover_servers(void)
{
if(is_server(i))
{
- sprintf(str_addr, "ssm://tcp::%s|%d:", ETHERNET_LOOPBACK, PORT_BASE+i);
+ sprintf(str_addr, "ssm://tcp::%s|%d:", str_ip_addresses[i].ip, port_base+i);
tret = triton_addr_lookup(str_addr, &server_addrs[nb_servers++]);
triton_assert(tret == TRITON_SUCCESS);
}
@@ -1038,7 +1085,7 @@ static void print_test_description()
{
if(rank != 0)
return;
- printf("\n\nThe program requires a unique string argument which is the input file name\n"
+ printf("\n\n Usage is: prog_name input_file_name [IP_Addresses_file_name]\n"
"See README_FOR_ONE_SIDED in \x1b[34mtriton/code/src/net/tests/\x1b[0m.\n"
"A sample input file named sample_one_sided_test_input can be copied from the aforementioned location as well\n\n");
}
@@ -1051,13 +1098,7 @@ __blocking int entry_point(int argc, char** argv)
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
MPI_Comm_size(MPI_COMM_WORLD, &size);
- set_environment();
-
-
- JZ_WAIT_FOR_DEBUGGER();
- tret = triton_msg_ssm_init();
- triton_assert(tret == TRITON_SUCCESS);
- if(argc != 2)
+ if(argc < 2)
{
if(rank == 0)
{
@@ -1077,6 +1118,19 @@ __blocking int entry_point(int argc, char** argv)
}
return 0;
}
+ if(argc == 3)
+ strcpy(ip_file_name, argv[2]);
+ JZ_WAIT_FOR_DEBUGGER();
+ if(build_ips() != TRITON_SUCCESS)
+ {
+ PRINT("Aborting for failure to get IP addreses");
+ MPI_Abort(MPI_COMM_WORLD, 0);
+ }
+ set_environment();
+
+
+ tret = triton_msg_ssm_init();
+ triton_assert(tret == TRITON_SUCCESS);
strcpy(input_file_name, argv[1]);
open_endpoints();
open_input();
diff --git a/code/src/net/tests/ssm-sr-test.ae b/code/src/net/tests/ssm-sr-test.ae
index 5a09386..7cfb37e 100644
--- a/code/src/net/tests/ssm-sr-test.ae
+++ b/code/src/net/tests/ssm-sr-test.ae
@@ -30,7 +30,6 @@ static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, c
}
#ifdef JZ_DEBUG
- void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once);
#define JZ_WAIT_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0)
#define JZ_WAIT_FOR_DEBUGGER_IF(cond) \
do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0);}while(0)
@@ -75,7 +74,7 @@ do
*/
-#define PORT_BASE 5000
+#define DEFAULT_PORT_BASE 5000
#define MAX_DATA_SIZE 50*MB
int rank = -1;
@@ -88,6 +87,7 @@ triton_msg_group_t group;
triton_msg_tag_t tag;
triton_addr_t server_addrs[MAX_NB_SERVERS];
+uint16_t port_base = DEFAULT_PORT_BASE;
static int is_prime(int n)
{
@@ -124,7 +124,7 @@ static void set_environment()
sprintf(env_string, "%s=%s", ENV_ETHERNET_IP, ETHERNET_LOOPBACK);
ret = putenv(env_string);
triton_assert(ret == 0);
- sprintf(env_string, "%s=%d", ENV_ETHERNET_PORT, PORT_BASE+rank);
+ sprintf(env_string, "%s=%d", ENV_ETHERNET_PORT, port_base+rank);
ret = putenv(env_string);
triton_assert(ret == 0);
}
@@ -162,7 +162,7 @@ static __blocking void discover_servers(void)
{
if(is_server(i))
{
- sprintf(str_addr, "ssm://tcp::%s|%d", ETHERNET_LOOPBACK, PORT_BASE+i);
+ sprintf(str_addr, "ssm://tcp::%s|%d", ETHERNET_LOOPBACK, port_base+i);
tret = triton_addr_lookup(str_addr, &server_addrs[nb_servers++]);
triton_assert(tret == TRITON_SUCCESS);
}
@@ -190,6 +190,12 @@ __blocking int entry_point(int argc, char** argv)
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
MPI_Comm_size(MPI_COMM_WORLD, &size);
+ if(argc ==2)
+ {
+ i = atoi(argv[1]);
+ if(i)
+ port_base = (uint16_t)i;
+ }
set_environment();
@@ -208,6 +214,7 @@ __blocking int entry_point(int argc, char** argv)
triton_assert(tret == TRITON_SUCCESS)
for(i=0; i<DATA_SIZE; i++)
buffer[i] = rank+100;
+ sleep(2);
tret = triton_msg_send(server_addrs[0], group, tag, 1, &buf, &data_size);
triton_assert(tret != TRITON_SUCCESS)
}
@@ -216,14 +223,22 @@ __blocking int entry_point(int argc, char** argv)
tret = triton_msg_recv_any(group, &requester_addr, &unused_tag, 1,
&buf, &max_recv_size, &received_bytes);
triton_assert(tret == TRITON_SUCCESS);
+ TRACE("requester_addr = (l=%lx, u=%lx)\n", requester_addr.l, requester_addr.u);
+ TRACE("unused_tag = %u\n", unused_tag);
for(i=0; i<DATA_SIZE; i++)
printf("buffer[%d] = %d\n", i, buffer[i]);
- tret = triton_msg_recv(requester_addr, group, unused_tag, 1,
+ /*
+ tret = triton_msg_recv_any(group, &requester_addr, &unused_tag, 1,
+ &buf, &max_recv_size, &received_bytes);
+ triton_assert(tret == TRITON_SUCCESS);*/
+ tret = triton_msg_recv(requester_addr, group, unused_tag+1, 1,
&buf, &max_recv_size, &received_bytes);
triton_assert(tret == TRITON_SUCCESS);
for(i=0; i<DATA_SIZE; i++)
printf("buffer[%d] = %d\n", i, buffer[i]);
}
+ MPI_Barrier(MPI_COMM_WORLD);
+ //JZ_WAIT_FOR_DEBUGGER();
close_endpoints();
triton_msg_ssm_finalize();
hooks/post-receive
--
1
0
28 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via 340ce550471845578fb757c6e55931cc0f1bd486 (commit)
from f11a0cbcfe05751011104fba685afc2d44ad5f85 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit 340ce550471845578fb757c6e55931cc0f1bd486
Author: Judicael Zounmevo <zounmevoj(a)mcs.anl.gov>
Date: Tue Aug 28 17:08:03 2012 -0500
Updated net-modules and 1st SSM net-module commit:
-Updated the net-interface and the triton net-layer
(renamed functions, removed wait_handles, and support for
variable mem_pub_handle serialization buffer size)
-Removed the group parameter from one-sided routines in
the triton net-layer
-Updated the MPI, Fault-injector and Mock net-modules
-Added the SSM net-module
IMPORTANT: Still waiting for bug fix in SSM (Shane) to have the net-module
fully functional
-Updated configure.ac for ssm dependencies and net-module
-Updated triton-error.spec for new needs
-Added a complete typical use for the SSM net-module
-Added a two-sided test for the ssm net-module
-Updated the sample input file format to make it net-module-
agnostic. The existing MPI test is updated consequently
-Exposed the triton_sem_t type in header
-----------------------------------------------------------------------
Summary of changes:
code/configure.ac | 30 +
code/src/common/resources/scheduling/sem.c | 42 +-
code/src/common/resources/scheduling/sem.hae | 10 +-
code/src/common/triton-error.spec | 5 +
code/src/net/fault-injector/fault-method.ae | 9 +-
code/src/net/mock/mock-method.ae | 8 +-
code/src/net/mpi/mpi-method.ae | 25 +-
.../triton-message-method.h | 5 +-
code/src/net/ssm/module.mk.in | 10 +
code/src/net/ssm/ssm-method.ae | 1634 ++++++++++++++++++++
code/src/net/ssm/ssm-method.h | 8 +
code/src/net/ssm/ssm-tcp-transport.c | 161 ++
code/src/net/ssm/ssm-transport.h | 69 +
code/src/net/tests/module.mk.in | 4 +-
code/src/net/tests/mpi-one-sided.ae | 101 +-
code/src/net/tests/sample_one_sided_test_input | 8 +-
.../tests/{mpi-one-sided.ae => ssm-one-sided.ae} | 210 ++-
code/src/net/tests/ssm-sr-test.ae | 233 +++
code/src/net/triton-message-method.hae | 10 +-
code/src/net/triton-message.ae | 135 +-
code/src/net/triton-message.hae | 61 +-
21 files changed, 2538 insertions(+), 240 deletions(-)
create mode 100644 code/src/net/ssm/module.mk.in
create mode 100644 code/src/net/ssm/ssm-method.ae
create mode 100644 code/src/net/ssm/ssm-method.h
create mode 100644 code/src/net/ssm/ssm-tcp-transport.c
create mode 100644 code/src/net/ssm/ssm-transport.h
copy code/src/net/tests/{mpi-one-sided.ae => ssm-one-sided.ae} (83%)
create mode 100644 code/src/net/tests/ssm-sr-test.ae
Diff of changes:
diff --git a/code/configure.ac b/code/configure.ac
index 0933f1e..fa5e240 100644
--- a/code/configure.ac
+++ b/code/configure.ac
@@ -176,6 +176,34 @@ fi
dnl todo: add check making sure that if the user manually enabled perftools
dnl and we cannot find them we complain
+
+
+dnl =========================================================================
+dnl == SSM ===================================================
+dnl =========================================================================
+
+AC_ARG_WITH(ssm,
+ AS_HELP_STRING([--with-ssm=dir],
+ [Location of ssm installation]),
+ SSM_DIR="$withval",SSM_DIR="")
+
+echo "SSM_DIR=${SSM_DIR}"
+
+#if test -n "$SSM_DIR" ;
+#then
+ # echo "Setting flags for ssm"
+ CPPFLAGS="$CPPFLAGS -I${SSM_DIR}/include"
+ LDFLAGS="$LDFLAGS -L${SSM_DIR}/lib"
+#fi
+
+
+AC_CHECK_HEADERS([ssm.h],[],
+ [AC_MSG_ERROR("ssm.h not found. Please specify a valid --with-ssm")])
+
+AC_CHECK_LIB([ssm],[ssm_stop],[],[AC_MSG_ERROR("libssm not found. Please specify a valid --with-ssm")],[-lpthread -lssmptcp])
+AC_CHECK_LIB([ssmptcp],[ssmptcp_new_tp],[],[AC_MSG_ERROR("libssmptcp not found. Please specify a valid --with-ssm")],[-lssm -lpthread])
+
+
dnl ======================================================================
dnl Look for MPI; Auto-enable if possible, unless disabled by the user
dnl ======================================================================
@@ -428,6 +456,8 @@ doc/resilience/module.mk
doc/resilience/resilience-book.txt
doc/resilience/prototype-2012-07/module.mk
src/asg/module.mk
+src/net/ssm/module.mk
src/examples/parallel-histogram/module.mk
src/examples/block-read-modify-write/module.mk
])
+
diff --git a/code/src/common/resources/scheduling/sem.c b/code/src/common/resources/scheduling/sem.c
index 9f8e1a2..78b3fd8 100644
--- a/code/src/common/resources/scheduling/sem.c
+++ b/code/src/common/resources/scheduling/sem.c
@@ -3,12 +3,6 @@
#include <assert.h>
-struct triton_sem
-{
- triton_mutex_t lock;
- unsigned int count;
- ae_ops_t queue;
-};
static int initialized = 0;
@@ -16,7 +10,6 @@ static ae_opcache_t triton_sem_opcache;
static int triton_sem_resource_id;
-
static inline void check_init (void)
{
assert (initialized && "sem resource called without being initialized??");
@@ -163,6 +156,7 @@ static struct ae_resource triton_sem_resource =
};
+
static triton_ret_t triton_sem_resource_init (void)
{
assert (!initialized);
@@ -187,6 +181,40 @@ static void triton_sem_resource_finalize (void)
ae_opcache_destroy (triton_sem_opcache);
}
+static int init_count = 0;
+static int is_thread_safe_init = 0; /*If thread-safe init is used once,
+ triton_sem_resource_init can no longer
+ be called directly in the app. Every other call must
+ use th etrhead safe version. The call must also
+ be freed with the thread-safe version
+ */
+
+static int is_indirect_init = 0; //this is set while a thread safe init is in progress
+
+/*Provided to enable reference counted-initialization
+*/
+triton_ret_t triton_sem_init_thread_safe(triton_mutex_t *mutex)
+{
+ is_thread_safe_init = 1;
+ triton_mutex_lock(mutex);
+ is_indirect_init = 1;
+ if(++init_count == 1)
+ triton_sem_resource_init();
+ is_indirect_init = 0;
+ triton_mutex_unlock(mutex);
+}
+
+triton_ret_t triton_sem_finalize_thread_safe(triton_mutex_t *mutex)
+{
+ triton_mutex_lock(mutex);
+ assert(init_count >= 0);
+ is_indirect_init = 1;
+ if(--init_count == 0)
+ triton_sem_resource_finalize();
+ is_indirect_init = 0;
+ triton_mutex_unlock(mutex);
+}
+
__attribute__((constructor)) void triton_sem_resource_init_register(void);
__attribute__((constructor)) void triton_sem_resource_init_register(void)
diff --git a/code/src/common/resources/scheduling/sem.hae b/code/src/common/resources/scheduling/sem.hae
index 7795181..d2d467c 100644
--- a/code/src/common/resources/scheduling/sem.hae
+++ b/code/src/common/resources/scheduling/sem.hae
@@ -10,8 +10,12 @@
* Only supports a count of 1 for up and down right now.
*/
-struct triton_sem;
-typedef struct triton_sem triton_sem_t;
+typedef struct triton_sem
+{
+ triton_mutex_t lock;
+ unsigned int count;
+ ae_ops_t queue;
+}triton_sem_t;
/**
* Initialize the semaphore. Sets internal count to the given value.
@@ -37,5 +41,7 @@ triton_ret_t triton_sem_up (triton_sem_t * sem);
*/
__blocking triton_ret_t triton_sem_down (triton_sem_t * sem);
+triton_ret_t triton_sem_init_thread_safe(triton_mutex_t *mutex);
+triton_ret_t triton_sem_finalize_thread_safe(triton_mutex_t *mutex);
#endif
diff --git a/code/src/common/triton-error.spec b/code/src/common/triton-error.spec
index 0ab560e..c3c94ba 100644
--- a/code/src/common/triton-error.spec
+++ b/code/src/common/triton-error.spec
@@ -42,6 +42,7 @@ ERR_DEADLOCK "Deadlock"
ERR_SHORT_IO "Unexpected short I/O operation"
ERR_IO "I/O error"
ERR_NO_TXN "Operation attempted on invalid or closed transaction"
+ERR_TXN_ALREADY "Transaction is already open"
ERR_NODE_INVALID "Invalid node"
ERR_RECV_TOO_SMALL "Receive buffer too small for message"
ERR_HINT_MISSING_TYPE "Hint type has not been registered"
@@ -61,3 +62,7 @@ ERR_MEM_IN_USE "Rejected attempt to unregister a memory with pu
ERR_MEM_OUT_OF_RANGE "Specified memory region is out of range"
ERR_NOT_IMPLEMENTED "The requested functionality is not (yet) implemented"
ERR_CONDITIONAL "Operation failed conditional check"
+ERR_MEM_REGISTERED "The memory is already registered"
+ERR_SSM_UNKNOWN "Unknown SSM error"
+ERR_SSM_TRSPT_TYPE_UNAVAILABLE "SSM transport type unavailable at this endpoint"
+ERR_NAME_HASH_COLLISION "Name hash collision"
diff --git a/code/src/net/fault-injector/fault-method.ae b/code/src/net/fault-injector/fault-method.ae
index 5326300..3723844 100644
--- a/code/src/net/fault-injector/fault-method.ae
+++ b/code/src/net/fault-injector/fault-method.ae
@@ -294,8 +294,7 @@ static __blocking triton_ret_t unpublish_memory(triton_mem_pub_handle_t handle /
return TRITON_ERR_NOT_IMPLEMENTED;
}
-
-static int get_serialized_handle_size(triton_serialized_handle_type_t handle_type)
+static int get_serialized_handle_size(triton_serialized_handle_type_t handle_type, void* handle)
{
return -1;
}
@@ -307,7 +306,7 @@ static triton_ret_t get_serialized_handle(triton_serialized_handle_type_t handle
return TRITON_ERR_NOT_IMPLEMENTED;
}
-static triton_ret_t set_serialized_handle(triton_serialized_handle_type_t handle_type,
+static triton_ret_t create_handle_from_serialization(triton_serialized_handle_type_t handle_type,
void* pointer_to_out_handle,
char* serialized_handle)
{
@@ -351,7 +350,7 @@ static __blocking triton_ret_t end_epoch(triton_mem_pub_handle_t handle)
}
static __blocking triton_ret_t wait(triton_mem_pub_handle_t mem_handle,
- triton_wait_handle_t wait_handle)
+ void* wait_info)
{
return TRITON_ERR_NOT_IMPLEMENTED;
}
@@ -382,7 +381,7 @@ static struct triton_msg_method triton_msg_fault_method =
.publish_memory = publish_memory,
.unpublish_memory = unpublish_memory,
.get_serialized_handle_size = get_serialized_handle_size,
- .set_serialized_handle = set_serialized_handle,
+ .create_handle_from_serialization = create_handle_from_serialization,
.get_serialized_handle = get_serialized_handle,
.free_serialized_handle = free_serialized_handle,
.get = get,
diff --git a/code/src/net/mock/mock-method.ae b/code/src/net/mock/mock-method.ae
index 636186d..2b08aa5 100644
--- a/code/src/net/mock/mock-method.ae
+++ b/code/src/net/mock/mock-method.ae
@@ -935,7 +935,7 @@ static __blocking triton_ret_t unpublish_memory(triton_mem_pub_handle_t handle /
}
-static int get_serialized_handle_size(triton_serialized_handle_type_t handle_type)
+static int get_serialized_handle_size(triton_serialized_handle_type_t handle_type, void* handle)
{
return -1;
}
@@ -947,7 +947,7 @@ static triton_ret_t get_serialized_handle(triton_serialized_handle_type_t handle
return TRITON_ERR_NOT_IMPLEMENTED;
}
-static triton_ret_t set_serialized_handle(triton_serialized_handle_type_t handle_type,
+static triton_ret_t create_handle_from_serialization(triton_serialized_handle_type_t handle_type,
void* pointer_to_out_handle,
char* serialized_handle)
{
@@ -991,7 +991,7 @@ static __blocking triton_ret_t end_epoch(triton_mem_pub_handle_t handle)
}
static __blocking triton_ret_t wait(triton_mem_pub_handle_t mem_handle,
- triton_wait_handle_t wait_handle)
+ void* wait_info)
{
return TRITON_ERR_NOT_IMPLEMENTED;
}
@@ -1022,7 +1022,7 @@ static struct triton_msg_method triton_msg_mock_method =
.publish_memory = publish_memory,
.unpublish_memory = unpublish_memory,
.get_serialized_handle_size = get_serialized_handle_size,
- .set_serialized_handle = set_serialized_handle,
+ .create_handle_from_serialization = create_handle_from_serialization,
.get_serialized_handle = get_serialized_handle,
.free_serialized_handle = free_serialized_handle,
.get = get,
diff --git a/code/src/net/mpi/mpi-method.ae b/code/src/net/mpi/mpi-method.ae
index 5bed0bf..c2a32c4 100644
--- a/code/src/net/mpi/mpi-method.ae
+++ b/code/src/net/mpi/mpi-method.ae
@@ -849,13 +849,13 @@ static __blocking triton_ret_t publish_memory(triton_mem_reg_handle_t in_handle,
if(!mem_pub_info)
return TRITON_ERR_NOMEM;
- triton_mutex_lock(&global_data.mutex);
mem_pub_info->offset = offset + ((uint64_t)mem_reg_info->base_address -
(uint64_t)global_data.one_sided_memory);
mem_pub_info->size = size;
mem_pub_info->mem_reg_info = mem_reg_info;
mem_pub_info->client_rank = global_data.rank;
+ triton_mutex_lock(&global_data.mutex);
mem_reg_info->publish_count++;
triton_mutex_unlock(&global_data.mutex);
@@ -911,14 +911,12 @@ static __blocking triton_ret_t unpublish_memory(triton_mem_pub_handle_t handle /
}
-static int get_serialized_handle_size(triton_serialized_handle_type_t handle_type)
+static int get_serialized_handle_size(triton_serialized_handle_type_t handle_type, void* handle)
{
switch(handle_type)
{
case MEM_PUB_HANDLE_TYPE:
return (int)sizeof(mem_pub_t);
- case WAIT_HANDLE_TYPE:
- return (int) sizeof(triton_wait_handle_t);
default:
triton_assert("unknown triton_serialized_handle_type_t");
}
@@ -946,23 +944,17 @@ static triton_ret_t get_serialized_handle(triton_serialized_handle_type_t handle
pi = (int*)((char*)pu64 + sizeof(uint64_t) + sizeof(int));//client_rank
*pi = aehton32(*pi);
break;
- case WAIT_HANDLE_TYPE:
- memcpy(serialized_handle, *((triton_wait_handle_t**)pointer_to_in_handle),
- sizeof(triton_wait_handle_t));
- *((uint64_t*)serialized_handle) = aehton64(*((uint64_t*)serialized_handle));
- break;
default:
return TRITON_ERR_INVAL;
}
return TRITON_SUCCESS;
}
-static triton_ret_t set_serialized_handle(triton_serialized_handle_type_t handle_type,
+static triton_ret_t create_handle_from_serialization(triton_serialized_handle_type_t handle_type,
void* pointer_to_out_handle,
char* serialized_handle)
{
mem_pub_t* mem_pub = NULL;
- uint64_t u64;
switch(handle_type)
{
case MEM_PUB_HANDLE_TYPE:
@@ -976,10 +968,6 @@ static triton_ret_t set_serialized_handle(triton_serialized_handle_type_t handle
mem_pub->made_from_deserialization = 1;
*((triton_mem_pub_handle_t*)pointer_to_out_handle) = (triton_mem_pub_handle_t)mem_pub;
break;
- case WAIT_HANDLE_TYPE:
- memcpy(&u64, serialized_handle, sizeof(triton_wait_handle_t));
- *((triton_wait_handle_t*)pointer_to_out_handle) = (triton_wait_handle_t)aentoh64(u64);
- break;
default:
return TRITON_ERR_INVAL;
}
@@ -1000,9 +988,6 @@ static triton_ret_t free_serialized_handle(triton_serialized_handle_type_t handl
free(mem_pub);
*((triton_mem_pub_handle_t*)pointer_to_in_handle) = 0;
break;
- case WAIT_HANDLE_TYPE:
- *((triton_wait_handle_t*)pointer_to_in_handle) = 0;
- break;
default:
return TRITON_ERR_INVAL;
}
@@ -1209,7 +1194,7 @@ static __blocking triton_ret_t end_epoch(triton_mem_pub_handle_t handle)
static __blocking triton_ret_t wait(triton_mem_pub_handle_t mem_handle,
- triton_wait_handle_t wait_handle)
+ void* wait_info)
{
//Noop. . . For this net-module, the orchestration of the
//completion guarantee is 100% managed in the upper layer
@@ -1241,7 +1226,7 @@ static struct triton_msg_method triton_msg_mpi_method =
.publish_memory = publish_memory,
.unpublish_memory = unpublish_memory,
.get_serialized_handle_size = get_serialized_handle_size,
- .set_serialized_handle = set_serialized_handle,
+ .create_handle_from_serialization = create_handle_from_serialization,
.get_serialized_handle = get_serialized_handle,
.free_serialized_handle = free_serialized_handle,
.get = get,
diff --git a/code/src/net/mpi/one_sided_interface_mpiimpl_2/triton-message-method.h b/code/src/net/mpi/one_sided_interface_mpiimpl_2/triton-message-method.h
index 5d4661c..58c82c4 100644
--- a/code/src/net/mpi/one_sided_interface_mpiimpl_2/triton-message-method.h
+++ b/code/src/net/mpi/one_sided_interface_mpiimpl_2/triton-message-method.h
@@ -242,7 +242,10 @@ struct triton_msg_method
__blocking triton_ret_t (*end_epoch)(triton_mem_pub_handle_t in_handle);
- __blocking triton_ret_t (*wait)(triton_mem_pub_handle_t handle, triton_wait_handle_t wait_handle);
+ /*
+ wait_context is a net-module specific data
+ */
+ __blocking triton_ret_t (*wait)(triton_mem_pub_handle_t handle, void* wait_context);
};
diff --git a/code/src/net/ssm/module.mk.in b/code/src/net/ssm/module.mk.in
new file mode 100644
index 0000000..60993a1
--- /dev/null
+++ b/code/src/net/ssm/module.mk.in
@@ -0,0 +1,10 @@
+DIR := src/net/ssm
+
+LIBSRC += $(DIR)/ssm-tcp-transport.c
+
+AELIBSRC += $(DIR)/ssm-method.ae
+
+#MODCFLAGS_$(DIR)/mpi-method = $(MPICFLAGS)
+#MODCC_$(DIR)/mpi-method = $(MPICC)
+
+#endif # BUILD_MPI
diff --git a/code/src/net/ssm/ssm-method.ae b/code/src/net/ssm/ssm-method.ae
new file mode 100644
index 0000000..452b7cb
--- /dev/null
+++ b/code/src/net/ssm/ssm-method.ae
@@ -0,0 +1,1634 @@
+/*Update the semaphores and mutex to use the posix ones if we can reason about the length of the resources they protect and that length is short. In particular, if they are protecting CPU-bound stuff, then, they might be plain sems and mutexes
+*/
+
+#include "ssm-method.h"
+#include "ssm-transport.h"
+#include "ssm.h"
+#include "ssmptcp.h"
+
+#include <pthread.h>
+#include <getopt.h>
+#include <assert.h>
+#include <unistd.h>
+#include <stdio.h>
+#include <limits.h>
+
+#include <sys/time.h>
+
+
+#include "src/common/triton-hash.h"
+
+#include "src/net/triton-addr.h"
+#include "src/net/triton-message.hae"
+#include "src/net/triton-message-method.hae"
+#include "src/common/triton-init.h"
+#include "src/remote/byteswap.h"
+#include "src/common/resources/scheduling/sem.hae"
+
+#include <endian.h>
+
+#include <semaphore.h>
+
+
+
+#define JZ_DEBUG
+static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once)
+{
+ static int count = 0;
+
+ if(count && once)
+ return;
+ count++;
+ printf("process ( pid = %d | rank = %d ) is waiting in %s at %s:%d for debugger\n",
+ getpid(), rank, function, file, line);
+ fflush(stdout);
+ for(;;);
+}
+
+#ifdef JZ_DEBUG
+ void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once);
+ #define JZ_WAIT_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0)
+ #define JZ_WAIT_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0);}while(0)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1);}while(0)
+#else
+ #define JZ_WAIT_FOR_DEBUGGER()
+ #define JZ_WAIT_FOR_DEBUGGER_IF()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF()
+#endif
+
+
+#define PRINT(...) do{printf(__VA_ARGS__); fflush(stdout);}while(0)
+#define TRACE(...) do{PRINT("%s:%d[%s]", __FILE__, __LINE__, __FUNCTION__); PRINT(__VA_ARGS__);}while(0)
+#define TRACE_IF(cond,...) do{if(cond)TRACE(__VA_ARGS__);}while(0)
+
+#define KB 1024
+#define MB 1048576
+
+typedef enum {FALSE, TRUE} BOOL;
+
+static __blocking triton_ret_t triton_sem_count_down_generic (triton_sem_t *sem, int count, BOOL ignore_cancel)
+{
+ int i;
+ triton_ret_t tret = TRITON_SUCCESS;
+ for (i=0; i<count; ++i)
+ {
+ tret = triton_sem_down(sem);
+ if(tret == TRITON_ERR_CANCELED)
+ {
+ if(!ignore_cancel)
+ break;
+ --i; //make up for the cancellation and wait anyways for the right count
+ }
+ }
+ return tret;
+}
+
+static inline __blocking triton_ret_t triton_sem_count_down (triton_sem_t *sem, int count)
+{
+ return triton_sem_count_down_generic(sem, count, FALSE);
+}
+
+static inline __blocking triton_ret_t triton_sem_count_down_no_cancel (triton_sem_t *sem, int count)
+{
+ return triton_sem_count_down_generic(sem, count, TRUE);
+}
+
+#define ENV_ETHERNET_IP "ETHERNET_IP"
+#define ENV_ETHERNET_PORT "ETHERNET_PORT"
+
+#define MAX_RAW_ADDRESS_SIZE 80
+
+#define ETHERNET_LOOPBACK "127.0.0.1"
+
+#define SSM_WAIT_TIMEOUT 2000 //2ms
+
+/*THE MATCH BITS MANAGEMENT STRATEGY:
+ 1- The most significant bit distinguishes between 2-sided and one-sided. If set, it is 2-sided
+ 2- For 2-sided:
+ The group name is hashed over the 16 next most significant bits
+ The tags which can therefore range from 0 to 0x7fff cover the 15 remaining bits
+ 3- One-sided range from 0 to 0x7ffffffe
+ 4- 0x7fffffff is reserved for message progression thread shutdown. It is kind of used by
+ a process to message itself when it is in the shutdown process.
+ */
+
+
+#define TWO_SIDED_MATCH_BITS_FLAG 0x8000000000000000 //the most significant bit is always set for 2-sided
+#define TWO_SIDED_ANY_SOURCE_MATCH_MASK (0x07ffffff)
+#define TWO_SIDED_TAG_MASK 0x07ffffffLLU
+#define GROUP_MASK (~TWO_SIDED_TAG_MASK)
+#define IS_ANY_SOURCE_MATCH_BITS(match_bits) (((match_bits)&TWO_SIDED_ANY_SOURCE_MATCH_MASK)==0)
+
+#define MAX_ONE_SIDED_MATCH_BITS 0X7ffffffffffffffd
+#define PROGRESSION_SHUTDOWN_RESERVED_MATCH_BITS 0x7fffffffffffffff
+
+#define SSM_MAX_TWO_SIDED_SIZE (4*KB)
+#define SSM_MAX_INTERNAL_TWO_SIDED_SIZE (SSM_MAX_TWO_SIDED_SIZE+MAX_RAW_ADDRESS_SIZE)
+
+#define SSM_ANY_SOURCE 0
+/*
+ GENERIC NOTES:
+ 1- This net-module will potentially manage several transports
+ 2- TCP ssm_transport is being implemented now. The other transports will definitely be different on a few points.
+ */
+
+
+#define PRINT_ERR_MSG(err_msg_format,...) \
+ do \
+{ \
+ triton_err(triton_log_default, "%s:%d: [pid %d]" err_msg_format, \
+ __FILE__, __LINE__, getpid(), ##__VA_ARGS__); \
+}while(0)
+
+#define ABORT(err_msg_format,...) \
+ do \
+{ \
+ PRINT_ERR_MSG(err_msg_format,##__VA_ARGS__); \
+ exit(1); \
+}while(0)
+
+
+#define ABORT_IF(cond,err_msg_format,...) do{if((cond)) ABORT(err_msg_format,##__VA_ARGS__);}while(0)
+
+
+#define ZERO_SIZE_MEM_HANDLE 0x80000000CAFEBABE
+
+/*TRANSPORT_TYPE_FROM_ADDR is encoded in the triton_msg_method_addr_t as the most significant byte*/
+#define TRANSPORT_TYPE_MASK_FOR_ADDR 0xff00000000000000
+
+
+#define TRANSPORT_TYPE_FROM_ADDR(addr) ((((uint64_t)(addr))&TRANSPORT_TYPE_MASK_FOR_ADDR)>>56)
+
+
+static uint64_t ssm_one_sided_size_threshold = SSM_MAX_TWO_SIDED_SIZE;
+
+static struct triton_hash_table *group_table = NULL;
+static const char* transport_type_delimiter = "::";
+
+typedef struct network_endpoint_t
+{
+ ssm_Itp ssm_transport;
+ ssm_id ssm_instance;
+ ssm_Iaddr address_interface;
+ ssm_Haddr local_listener_address_handle;
+ ssm_me shutdown_me; //the progress engine shutdown is triggered by an ssm_unlink op.
+ unsigned char raw_address[MAX_RAW_ADDRESS_SIZE];
+ ssm_transport_abstraction_t* transport_abstraction;
+ ssm_transport_type_t transport_type;
+
+ sem_t release_sem; //used in the release phase
+
+ triton_list_link_t list_link;
+
+ struct triton_hash_link address_hash_link; //if the address is not encodable over 56-bit, this contains
+ //the actual address data. It uses the address as a key then
+}network_endpoint_t;
+
+network_endpoint_t tcp_network_endpoint;
+//TODO: Add network endpoints for other transports when they are available
+
+network_endpoint_t *default_net_endpoint;
+
+
+triton_msg_method_addr_t default_self_addr;
+
+
+struct
+{
+ triton_mutex_t mutex;
+ triton_mutex_t two_sided_mutex;
+ triton_list_t mem_reg_list;
+ triton_list_t net_endpoint_list;
+ triton_list_t generic_recv_data_list;
+
+ uint64_t one_sided_match_bits_counter; //ensure match bit unicity
+
+ triton_list_link_t ptp_match_entries_list;
+
+ pthread_t progress_thread;
+ BOOL exiting; //TRUE when the process is shuting down
+
+}ssm_global_data;
+
+
+
+static void shutdown_callback_routine(void* cb_args, void* ev_data)
+{
+ if(((ssm_result)ev_data)->op == SSM_OP_UNLINK)
+ {
+ sem_post(&((network_endpoint_t*)cb_args)->release_sem);
+ }
+}
+
+ssm_cb_t shutdown_callback;
+
+static void init_ethernet()
+{
+ char* env_val = getenv(ENV_ETHERNET_IP);
+ char env_val_buf[101];
+ unsigned short port;
+ int i = 0, ip_chunk;
+ char *token = NULL;
+
+ strcpy(env_val_buf, env_val? env_val:ETHERNET_LOOPBACK);
+ token = strtok(env_val_buf, ".");
+
+ while(token)
+ {
+ ip_chunk = atoi(token);
+ tcp_network_endpoint.raw_address[i++] = (unsigned char)ip_chunk;
+ token = strtok(NULL, ".");
+ ABORT_IF(i>4, "Wrong IPv4 address format specified in %s", ENV_ETHERNET_IP);
+ }
+ env_val = getenv(ENV_ETHERNET_PORT);
+ ABORT_IF(!env_val, "No Ethernet port specified to build local ethernet network endpoint");
+
+ port = (unsigned short)atoi(env_val);
+
+ memcpy(&tcp_network_endpoint.raw_address[4], &port, sizeof(unsigned short));
+
+ tcp_network_endpoint.ssm_transport = ssmptcp_new_tp(port, SSM_NOF);
+ tcp_network_endpoint.ssm_instance = ssm_start(tcp_network_endpoint.ssm_transport, NULL, SSM_NOF);
+ tcp_network_endpoint.address_interface = ssm_addr(tcp_network_endpoint.ssm_instance);
+ tcp_network_endpoint.local_listener_address_handle =
+ tcp_network_endpoint.address_interface->local(tcp_network_endpoint.address_interface);
+ tcp_network_endpoint.transport_abstraction = get_tcp_abstraction();
+ tcp_network_endpoint.transport_type = SSM_TCP;
+
+ shutdown_callback.pcb = shutdown_callback_routine;
+ shutdown_callback.cbdata = &tcp_network_endpoint;
+ tcp_network_endpoint.shutdown_me = ssm_link(tcp_network_endpoint.ssm_instance,
+ PROGRESSION_SHUTDOWN_RESERVED_MATCH_BITS, 0, SSM_POS_HEAD, NULL, &shutdown_callback, SSM_NOF);
+
+ sem_init(&tcp_network_endpoint.release_sem, 0, 0);
+
+ triton_list_link_clear(&tcp_network_endpoint.list_link);
+ triton_list_add_back(&tcp_network_endpoint.list_link, &ssm_global_data.net_endpoint_list);
+}
+/*Provide the init_ethernet equivalent for the other transports as well
+ */
+
+static inline void release_ethernet()
+{
+ ssm_unlink(tcp_network_endpoint.ssm_instance, tcp_network_endpoint.shutdown_me);
+ sem_wait(&tcp_network_endpoint.release_sem);
+ ssm_stop(tcp_network_endpoint.ssm_instance);
+ sem_destroy(&tcp_network_endpoint.release_sem);
+}
+
+
+static inline void set_default_net_endpoint()
+{
+ default_net_endpoint = &tcp_network_endpoint;
+}
+
+
+static ssm_transport_abstraction_t* get_transport_abstraction(ssm_transport_type_t type)
+{
+ switch(SSM_TCP)
+ {
+ case SSM_TCP:
+ return get_tcp_abstraction();
+ }
+ return NULL;
+}
+
+
+static inline ssm_transport_abstraction_t* get_transport_abstraction_from_addr(triton_msg_method_addr_t addr)
+{
+ return get_transport_abstraction(TRANSPORT_TYPE_FROM_ADDR(addr));
+}
+
+static network_endpoint_t* get_network_endpoint(ssm_transport_type_t type)
+{
+ switch(SSM_TCP)
+ {
+ case SSM_TCP:
+ return &tcp_network_endpoint;
+ }
+ return NULL;
+}
+
+/*
+Add a remote address that will be used afterward simply through a triton_msg_method_addr_t
+The address is not added if it already exists
+*/
+static void add_address(network_endpoint_t* net_endpoint, unsigned char* raw_host_bytes_addr_data, triton_msg_method_addr_t* addr)
+{
+ switch(net_endpoint->transport_type)
+ {
+ case SSM_TCP:
+ *addr = (triton_msg_method_addr_t)net_endpoint->transport_abstraction->get_address(
+ raw_host_bytes_addr_data, NULL);
+ break;
+
+ /*The other ones will probably have some kind of data structure that will will be created and keyed
+ with the triton_msg_method_addr_t
+ */
+ default:
+ ;
+ }
+}
+
+static __blocking triton_msg_method_addr_t ssm_addr_self(void)
+{
+ //TODO: How do we know what the caller of this function want if this ssm_instance of the net-module
+ //has several transports underneath (E.g. TCP, IB, UDP, some other stuff)
+ return (triton_msg_method_addr_t)default_net_endpoint->transport_abstraction->get_address(
+ default_net_endpoint->raw_address, NULL);
+}
+
+static int ssm_addr_equal(triton_msg_method_addr_t a1, triton_msg_method_addr_t a2)
+{
+ return a1 == a2;
+}
+
+static triton_ret_t ssm_addr_to_string(triton_msg_method_addr_t addr, triton_string_t *addrstr)
+{
+ char str_addr[MAX_RAW_ADDRESS_SIZE];
+ network_endpoint_t* net_endpoint = get_network_endpoint(TRANSPORT_TYPE_FROM_ADDR(addr));
+ if(!net_endpoint)
+ return TRITON_ERR_SSM_TRSPT_TYPE_UNAVAILABLE;
+ net_endpoint->transport_abstraction->get_str_address((uint64_t)addr, str_addr, net_endpoint->raw_address);
+ triton_string_init(addrstr, "%s", str_addr);
+ return TRITON_SUCCESS;
+}
+
+static triton_ret_t ssm_addr_free(triton_msg_method_addr_t addr)
+{
+ switch(TRANSPORT_TYPE_FROM_ADDR(addr))
+ {
+ case SSM_TCP:
+ break; //nothing to do for TCP; it uses the type|address enconding
+ }
+ //The other ones in the future might be reference-counted and freed only when all copies are freed
+ return TRITON_SUCCESS;
+}
+
+static triton_ret_t ssm_addr_copy(const triton_msg_method_addr_t orig_addr, triton_msg_method_addr_t *copy)
+{
+ switch(TRANSPORT_TYPE_FROM_ADDR(orig_addr))
+ {
+ case SSM_TCP:
+ *copy = orig_addr;
+ break; //Other ssm_transport that maintain additional data might need reference counting
+ }
+ return TRITON_SUCCESS;
+}
+
+triton_msg_method_addr_t* g_looked_up_addr;
+
+static triton_ret_t ssm_addr_lookup(const char *name, triton_msg_method_addr_t *addr)
+{
+ char name_1[251];
+ triton_assert(strlen(name)<=250);
+ if(strncmp(name, "tcp", 3) == 0)
+ {
+ strcpy(name_1, name);
+ *addr = (triton_msg_method_addr_t)tcp_network_endpoint.transport_abstraction->get_address_from_str_addr(name_1, NULL);
+ g_looked_up_addr = addr;
+ }
+ else
+ {
+ triton_assert("Unsupported ssm_transport!");
+ }
+ return TRITON_SUCCESS;
+}
+
+ssmptcp_addrargs_t* g_tcp_args;
+static ssm_Haddr tcp_get_ssm_addr_from_triton_addr(triton_msg_method_addr_t addr)
+{
+ unsigned char raw_address[6];
+ unsigned char address_arg[SSM_MAX_ADDR_ARG_SIZE];
+ ssmptcp_addrargs_t* tcp_args = NULL;
+ ssm_transport_abstraction_t * tcp_abstraction = get_tcp_abstraction();
+
+ uint64_t u64_addr = (uint64_t)addr;
+
+ raw_address[0] = (unsigned char)((u64_addr&0xff0000000000)>>40);
+ raw_address[1] = (unsigned char)((u64_addr&0xff00000000)>>32);
+ raw_address[2] = (unsigned char)((u64_addr&0xff000000)>>24);
+ raw_address[3] = (unsigned char)((u64_addr&0xff0000)>>16);
+
+ *((uint16_t*)(&raw_address[4])) = u64_addr&0xffff;
+
+ tcp_abstraction->create_address_arg_data(raw_address, address_arg);
+ tcp_args = (ssmptcp_addrargs_t*)address_arg;
+ g_tcp_args = tcp_args;
+ return tcp_abstraction->create_addr_handle(tcp_network_endpoint.address_interface, tcp_args);
+}
+
+static triton_ret_t get_ssm_addr_from_triton_addr(triton_msg_method_addr_t addr, ssm_Haddr* the_ssm_addr)
+{
+
+ switch(SSM_TCP)
+ {
+ case SSM_TCP:
+ *the_ssm_addr = tcp_get_ssm_addr_from_triton_addr(addr);
+ return TRITON_SUCCESS;
+ }
+ return TRITON_ERR_SSM_TRSPT_TYPE_UNAVAILABLE;
+}
+
+/*
+ */
+
+
+#define SSM_TAG_MAX TWO_SIDED_TAG_MASK
+struct triton_ssm_group
+{
+ char *name;
+ triton_msg_method_group_t group;
+ int tag_counter;
+ triton_mutex_t mutex;
+ struct triton_hash_link link;
+};
+
+
+static triton_ret_t ssm_group_add(const char *service_name, triton_msg_method_group_t *group)
+{
+ uint32_t h1 = 0, h2 = 0;
+ struct triton_ssm_group *newgroup;
+ BOOL group_exists_already = FALSE;
+ assert(service_name);
+ assert(group);
+
+ bj_hashlittle2(service_name, strlen(service_name), &h1, &h2);
+
+ triton_mutex_lock(&ssm_global_data.mutex);
+ if(triton_hash_search(group_table, (triton_msg_method_group_t*)(&h1)) != NULL)
+ group_exists_already = TRUE;
+ triton_mutex_unlock(&ssm_global_data.mutex);
+ if(group_exists_already)
+ return TRITON_ERR_NAME_HASH_COLLISION;
+
+ newgroup = malloc(sizeof(*newgroup));
+ if(newgroup == NULL) return TRITON_ERR_NOMEM;
+
+ newgroup->name = strdup(service_name);
+
+ newgroup->group = h1;
+ newgroup->tag_counter = 0;
+ triton_mutex_init(&newgroup->mutex, NULL);
+ *group = newgroup->group;
+
+ triton_list_link_clear(&newgroup->link);
+ triton_hash_add(group_table, &newgroup->group, &newgroup->link);
+ return TRITON_SUCCESS;
+}
+
+static triton_msg_tag_t ssm_new_tag(triton_msg_method_group_t group)
+{
+ struct triton_hash_link *groupl;
+ struct triton_ssm_group *group_entry;
+ triton_msg_tag_t tag;
+
+ /* lookup group */
+ groupl = triton_hash_search(group_table, &group);
+ group_entry = triton_hash_get_entry(groupl, struct triton_ssm_group, link);
+ triton_mutex_lock(&group_entry->mutex);
+ ++(group_entry->tag_counter);
+ if(group_entry->tag_counter >= SSM_TAG_MAX)
+ {
+ group_entry->tag_counter = 0;
+ }
+ tag = group_entry->tag_counter;
+ triton_mutex_unlock(&group_entry->mutex);
+ return tag;
+}
+
+
+typedef struct send_comm_info_t
+{
+ ssm_id ssm_instance;
+ ssm_tx transaction;
+ ssm_md md;
+ ssm_mr mr;
+ triton_sem_t semaphore;
+ ssm_status status;
+ char* intermediate_buffer;
+}send_comm_info_t;
+
+
+static void send_callback(void* cb_data, void* event_data)
+{
+ static int counter = 0;
+ send_comm_info_t* comm_info = (send_comm_info_t*)cb_data;
+ ssm_result result = (ssm_result)event_data;
+ if(result->op == SSM_ST_CANCEL || (result->status == SSM_ST_COMPLETE && result->op == SSM_OP_PUT))
+ {
+ comm_info->status = result->status;
+ triton_sem_up(&comm_info->semaphore);
+ }
+ TRACE("Inside send_callback no = %d (status, op, ssm_bits) = ( %d, %d, 0x%llx)\n", counter++, result->status, result->op, (uint64_t)result->bits);
+}
+
+send_comm_info_t* g_send_comm_info;
+ssm_bits g_match_bits;
+/*For now, 2-sided operations are limited to a single contiguous buffer*/
+static __blocking triton_ret_t ssm_send(
+ triton_msg_method_addr_t to,
+ triton_msg_method_group_t group,
+ triton_msg_tag_t tag,
+ int count,
+ char **buffers,
+ uint32_t *sizes)
+{
+ triton_ret_t tret;
+ int i;
+ network_endpoint_t* net_endpoint;
+ send_comm_info_t* comm_info = NULL;
+ ssm_Haddr the_ssm_addr = NULL;
+ ssm_cb_t callback;
+ callback.pcb = send_callback;
+
+ struct triton_hash_link *groupl;
+ struct triton_ssm_group *group_entry;
+ ssm_bits match_bits;
+
+ if(count>1)
+ return TRITON_ERR_NOT_IMPLEMENTED; //for now, only contiguous stuff is supported
+ if(sizes[0] > SSM_MAX_TWO_SIDED_SIZE)
+ return TRITON_ERR_MEM_OUT_OF_RANGE;
+
+ tret = get_ssm_addr_from_triton_addr(to, &the_ssm_addr);
+ if(tret != TRITON_SUCCESS)
+ return tret;
+ net_endpoint = get_network_endpoint(TRANSPORT_TYPE_FROM_ADDR(to));
+ triton_assert(net_endpoint); //should not fail if get_ssm_addr_from_triton_addr succeded.
+
+ groupl = triton_hash_search(group_table, &group);
+ group_entry = triton_hash_get_entry(groupl, struct triton_ssm_group, link);
+
+ /*Most significant bit (msb) is always set. Then we take the 4 msb of the group and finally the tag*/
+ match_bits = (TWO_SIDED_MATCH_BITS_FLAG | (((uint64_t)group_entry->group)<<31ULL)|tag);
+ g_match_bits = match_bits;
+
+ comm_info = malloc(sizeof(send_comm_info_t) + SSM_MAX_INTERNAL_TWO_SIDED_SIZE);
+ if(!comm_info)
+ return TRITON_ERR_NOMEM;
+
+ g_send_comm_info = comm_info;
+ comm_info->intermediate_buffer = ((char*)comm_info + sizeof(send_comm_info_t));
+
+ /*build and add the sender address; it will be shipped with the data
+ The first MAX_RAW_ADDRESS_SIZE bytes are reserved for the serialized sender address
+ */
+ net_endpoint->transport_abstraction->get_serialized_address(
+ net_endpoint->ssm_transport,
+ net_endpoint->raw_address, (uint8_t*)comm_info->intermediate_buffer);
+
+ memcpy(comm_info->intermediate_buffer + MAX_RAW_ADDRESS_SIZE, buffers[0], sizes[0]);
+
+ comm_info->ssm_instance = net_endpoint->ssm_instance;
+ callback.cbdata = comm_info;
+ triton_sem_init(&comm_info->semaphore, 0);
+
+ comm_info->mr = ssm_mr_create(NULL, comm_info->intermediate_buffer,
+ sizes[0] + MAX_RAW_ADDRESS_SIZE);
+
+ comm_info->md = ssm_md_add(NULL, 0, sizes[0] + MAX_RAW_ADDRESS_SIZE);
+ comm_info->transaction = ssm_put(net_endpoint->ssm_instance, the_ssm_addr,
+ comm_info->mr, comm_info->md, match_bits, &callback, SSM_NOF);
+
+ tret = triton_sem_down(&comm_info->semaphore);
+
+ if(tret == TRITON_ERR_CANCELED)
+ {
+ ssm_cancel(comm_info->ssm_instance, comm_info->transaction);
+ triton_sem_count_down_no_cancel(&comm_info->semaphore, 1);
+ }
+ else
+ {
+ if(comm_info->status != SSM_ST_COMPLETE)
+ tret = TRITON_ERR_AGAIN;
+ }
+ ssm_mr_destroy(comm_info->mr);
+ ssm_md_release(comm_info->md);
+ triton_sem_destroy(&comm_info->semaphore);
+ free(comm_info);
+ TRACE("Exiting send\n");
+ return tret;
+}
+
+/*Each of the data structure below is created for each receive*/
+typedef struct specific_recv_comm_info_t
+{
+ triton_msg_method_addr_t *sender; //this is a required info for matching
+ triton_msg_tag_t *tag;
+ char* buffer;
+ uint32_t size;
+ uint32_t* received_bytes;
+ triton_sem_t semaphore;
+ triton_list_link_t list_link;
+}specific_recv_comm_info_t;
+
+
+typedef struct generic_recv_comm_info_t
+{
+ network_endpoint_t* net_endpoint;
+ ssm_me match_entry;
+ ssm_bits match_bits;
+ BOOL is_receive_any;
+ ssm_status status;
+ triton_list_link_t list_link;
+ triton_list_t specific_recv_list;
+}generic_recv_comm_info_t;
+
+
+static int specific_recv_comm_info_predicate(struct triton_list_link *ll, void *the_key)
+{
+ triton_msg_method_addr_t key = (triton_msg_method_addr_t)the_key;
+ specific_recv_comm_info_t* comm_info = triton_list_get_entry(ll, specific_recv_comm_info_t, list_link);
+ return *comm_info->sender == SSM_ANY_SOURCE || *comm_info->sender == key;
+}
+
+typedef struct generic_recv_comm_key_t
+{
+ ssm_id ssm_instance;
+ ssm_bits match_bits;
+ BOOL is_receive_any;
+}generic_recv_comm_key_t;
+
+static int generic_recv_comm_info_predicate(struct triton_list_link *ll, void *the_key)
+{
+ generic_recv_comm_key_t* key = (generic_recv_comm_key_t*)the_key;
+ generic_recv_comm_info_t* comm_info = triton_list_get_entry(ll, generic_recv_comm_info_t, list_link);
+ if(key->is_receive_any)
+ return comm_info->is_receive_any && comm_info->net_endpoint->ssm_instance == key->ssm_instance &&
+ (comm_info->match_bits&GROUP_MASK) == (key->match_bits&GROUP_MASK);
+ return comm_info->net_endpoint->ssm_instance == key->ssm_instance &&
+ comm_info->match_bits == key->match_bits;
+}
+
+/*It is assumed that each call to this function is enclosed in two_sided_mutex
+*/
+static specific_recv_comm_info_t* remove_receive(generic_recv_comm_info_t* generic_comm_info,
+ triton_msg_method_addr_t sender)
+{
+ TRACE("Inside remove_receive\n");
+ specific_recv_comm_info_t* specific_comm_info = NULL;
+
+ triton_list_link_t *ll = triton_list_find(&generic_comm_info->specific_recv_list,
+ specific_recv_comm_info_predicate, (void*)sender);
+ if(!ll)
+ return NULL;
+ specific_comm_info = triton_list_get_entry(ll, specific_recv_comm_info_t, list_link);
+ triton_list_del(&specific_comm_info->list_link);
+ return specific_comm_info;
+}
+
+static inline void cancel_receive(generic_recv_comm_info_t* generic_comm_info,triton_msg_method_addr_t sender)
+{
+ (void)remove_receive(generic_comm_info, sender);
+}
+
+static void recv_callback(void* cb_data, void* event_data)
+{
+ static int counter = 0;
+ JZ_WAIT_FOR_DEBUGGER();
+ TRACE("Entering recv_callback no = %d \n", counter);
+ int ret;
+ generic_recv_comm_info_t* comm_info = (generic_recv_comm_info_t*)cb_data;
+ ssm_result result = (ssm_result)event_data;
+ ssm_mrinfo_t mr_info;
+ ssm_mr mr = result->mr;
+ triton_msg_method_addr_t sender;
+ specific_recv_comm_info_t* specific_receive_info = NULL;
+
+ triton_assert(result->op == SSM_OP_PUT);
+
+ ret = ssm_mr_getinfo(&mr_info, mr);
+ triton_assert(ret == 0);
+ ssm_mr_destroy(mr);
+
+ //extract sender address
+ comm_info->net_endpoint->transport_abstraction->get_raw_address_in_host_bytes(
+ comm_info->net_endpoint->ssm_transport, mr_info.base, mr_info.base); //we transform mr_info.base in-place
+ //so that the address could be referenced afterward with a triton_msg_method_addr_t
+ add_address(comm_info->net_endpoint, mr_info.base, &sender);
+
+ triton_mutex_lock(&ssm_global_data.two_sided_mutex);
+ specific_receive_info = remove_receive(comm_info, sender);
+ if(!specific_receive_info)//probably canceled previously
+ {
+ triton_mutex_unlock(&ssm_global_data.two_sided_mutex);
+ free(mr_info.base);
+ return;
+ }
+ *specific_receive_info->received_bytes = result->bytes - MAX_RAW_ADDRESS_SIZE;
+ *specific_receive_info->sender = sender;
+ *specific_receive_info->tag = result->bits & TWO_SIDED_TAG_MASK;
+ if(specific_receive_info->size <= *specific_receive_info->received_bytes)
+ memcpy(specific_receive_info->buffer, (char*)mr_info.base + MAX_RAW_ADDRESS_SIZE,
+ mr_info.span - MAX_RAW_ADDRESS_SIZE);
+
+ triton_sem_up(&specific_receive_info->semaphore);
+
+ free(mr_info.base);
+
+ if(comm_info->specific_recv_list.count == 0)
+ {
+ ssm_unlink(comm_info->net_endpoint->ssm_instance, comm_info->match_entry);
+ triton_list_del(&comm_info->list_link);
+ free(comm_info);
+ }
+ triton_mutex_unlock(&ssm_global_data.two_sided_mutex);
+ TRACE("Exiting recv_callback no = %d (status, op, ssm_bits) = ( %d, %d, 0x%llx)\n", counter++, result->status, result->op, (uint64_t)result->bits);
+}
+
+
+generic_recv_comm_info_t* g_recv_comm;
+/*
+sender and tag are in/out parameters
+*/
+static __blocking triton_ret_t generic_receive(triton_msg_method_group_t group, triton_msg_tag_t *tag,
+ triton_msg_method_addr_t* sender, char* user_buffer, uint32_t buf_size, uint32_t* received_bytes,
+ BOOL is_receive_any)
+{
+ triton_ret_t tret;
+
+ struct triton_hash_link *groupl;
+ struct triton_ssm_group *group_entry;
+ generic_recv_comm_key_t key;
+ triton_list_link_t* ll = NULL;
+ generic_recv_comm_info_t* comm_info = NULL;
+ ssm_mr mr;
+ specific_recv_comm_info_t specific_recv_comm_info;
+ ssm_cb_t callback;
+ ssm_bits match_bits;
+ network_endpoint_t* net_endpoint = NULL;
+ char *buffer = NULL;
+
+ specific_recv_comm_info.sender = sender;
+ specific_recv_comm_info.buffer = user_buffer;
+ specific_recv_comm_info.size = buf_size;
+ specific_recv_comm_info.received_bytes = received_bytes;
+ specific_recv_comm_info.tag = tag;
+ triton_sem_init(&specific_recv_comm_info.semaphore, 0);
+
+ net_endpoint = is_receive_any ? default_net_endpoint: get_network_endpoint(TRANSPORT_TYPE_FROM_ADDR(*sender));
+ triton_assert(net_endpoint); //should not fail if get_ssm_addr_from_triton_addr succeded.
+
+ groupl = triton_hash_search(group_table, &group);
+ group_entry = triton_hash_get_entry(groupl, struct triton_ssm_group, link);
+
+ match_bits = (TWO_SIDED_MATCH_BITS_FLAG | (((uint64_t)group_entry->group)<<31ULL));
+ if(!is_receive_any)
+ match_bits |= *tag;
+
+ callback.pcb = recv_callback;
+ buffer = malloc(SSM_MAX_INTERNAL_TWO_SIDED_SIZE);
+ if(!buffer)
+ return TRITON_ERR_NOMEM;
+
+ key.ssm_instance = net_endpoint->ssm_instance;
+ key.match_bits = match_bits;
+ key.is_receive_any = is_receive_any;
+
+ triton_mutex_lock(&ssm_global_data.two_sided_mutex);
+ ll = triton_list_find(&ssm_global_data.generic_recv_data_list, generic_recv_comm_info_predicate, &key);
+ if(ll)
+ comm_info = triton_list_get_entry(ll, generic_recv_comm_info_t, list_link);
+ else
+ {
+ comm_info = malloc(sizeof(generic_recv_comm_info_t));
+ if(!comm_info)
+ {
+ triton_mutex_unlock(&ssm_global_data.two_sided_mutex);
+ triton_sem_destroy(&specific_recv_comm_info.semaphore);
+ free(buffer);
+ return TRITON_ERR_NOMEM;
+ }
+ comm_info->net_endpoint = net_endpoint;
+ comm_info->match_bits = match_bits;
+ triton_list_init(&comm_info->specific_recv_list);
+ triton_list_link_clear(&comm_info->list_link);
+ triton_list_add_back(&comm_info->list_link, &ssm_global_data.generic_recv_data_list);
+ callback.cbdata = comm_info;
+ comm_info->is_receive_any = is_receive_any;
+ comm_info->match_entry = ssm_link(net_endpoint->ssm_instance, match_bits,
+ is_receive_any?TWO_SIDED_ANY_SOURCE_MATCH_MASK:0, SSM_POS_TAIL, NULL, &callback, SSM_NOF);
+ g_recv_comm = comm_info;
+ }
+
+ triton_list_link_clear(&specific_recv_comm_info.list_link);
+ triton_list_add_back(&specific_recv_comm_info.list_link, &comm_info->specific_recv_list);
+
+ mr = ssm_mr_create(NULL, buffer, SSM_MAX_INTERNAL_TWO_SIDED_SIZE);
+ ssm_post(net_endpoint->ssm_instance, comm_info->match_entry, mr, SSM_NOF);
+
+ triton_mutex_unlock(&ssm_global_data.two_sided_mutex);
+
+ tret = triton_sem_down(&specific_recv_comm_info.semaphore);
+ TRACE("After triton_sem_down in Generic receive\n");
+ if(tret == TRITON_ERR_CANCELED)
+ {
+ triton_mutex_lock(&ssm_global_data.two_sided_mutex);
+ cancel_receive(comm_info, *sender);
+ triton_mutex_unlock(&ssm_global_data.two_sided_mutex);
+ triton_sem_destroy(&specific_recv_comm_info.semaphore);
+ return tret;
+ }
+ if(*received_bytes > buf_size)
+ tret = TRITON_ERR_RECV_TOO_SMALL;
+ triton_sem_destroy(&specific_recv_comm_info.semaphore);
+ TRACE("Exiting Generic receive\n");
+ return tret;
+}
+
+
+static __blocking triton_ret_t ssm_recv(
+ triton_msg_method_addr_t from,
+ triton_msg_method_group_t group,
+ triton_msg_tag_t tag,
+ int count,
+ char **buffers,
+ uint32_t *sizes,
+ uint32_t *bytes_received)
+{
+ if(count>1)
+ return TRITON_ERR_NOT_IMPLEMENTED; //for now, only contiguous stuff is supported
+ return generic_receive(group, &tag, &from, buffers[0], sizes[0], bytes_received, FALSE);
+}
+
+
+static __blocking triton_ret_t ssm_recv_any(
+ triton_msg_method_group_t group,
+ triton_msg_method_addr_t *from,
+ triton_msg_tag_t *tag,
+ int count,
+ char **buffers,
+ uint32_t *sizes,
+ uint32_t *bytes_received)
+{
+ if(count>1)
+ return TRITON_ERR_NOT_IMPLEMENTED; //for now, only contiguous stuff is supported
+ *from = SSM_ANY_SOURCE;
+ return generic_receive(group, tag, from, buffers[0], sizes[0], bytes_received, TRUE);
+}
+
+static int ssm_group_compare(const void *key, struct triton_hash_link *hash_link)
+{
+ const triton_msg_method_group_t *group1 = (const triton_msg_method_group_t *)key;
+ struct triton_ssm_group *ssm_group = triton_hash_get_entry(hash_link, struct triton_ssm_group, link);
+ assert(key);
+ assert(hash_link);
+
+ return (ssm_group->group == *group1);
+}
+
+static void group_entry_free(void *x)
+{
+ struct triton_ssm_group *g;
+
+ g = (struct triton_ssm_group *)x;
+
+ triton_mutex_destroy(&g->mutex);
+ free(g->name);
+ free(g);
+}
+
+/*---------------One-sided --------------------*/
+
+static triton_ret_t set_one_sided_threshold(size_t size)
+{
+ triton_assert((int64_t)size >= 0);
+ ssm_one_sided_size_threshold = size;
+ return TRITON_SUCCESS;
+}
+
+
+static triton_ret_t get_one_sided_threshold(size_t *size)
+{
+ *size = (size_t)ssm_one_sided_size_threshold;
+ return TRITON_SUCCESS;
+}
+
+
+typedef struct mem_reg_t
+{
+ void* base_address;
+ uint64_t size;
+ int publish_count;
+ BOOL is_registered; //FALSE if this mem_reg_t is only linked to a buffer_allocate'd buffer
+ BOOL is_internally_allocated;
+ triton_list_link_t list_link;
+ ssm_transport_type_t transport_type; //the ssm_transport type this memory has been registered with
+}mem_reg_t;
+
+
+typedef struct transaction_t
+{
+ ssm_tx transaction;
+ ssm_md md;
+ ssm_mr mr;
+ triton_list_link_t list_link; //useful for one-sided
+}transaction_t;
+
+typedef struct mem_pub_t
+{
+ mem_reg_t* mem_reg_info;
+ void* base;
+ uint64_t size;
+ ssm_mr mr;
+ ssm_bits match_bits;
+ ssm_transport_type_t transport_type;
+ ssm_id transport_instance;
+ ssm_Haddr initial_owner_address; //This is supposed to be built on the side of the one-sided origin peer
+ ssm_cb_t callback;
+ int nb_expected_completions; //for the target
+ triton_mutex_t mutex;
+ triton_sem_t comm_semaphore; //MUST be used EXCLUSIVELY for communications. It is not a generic sem
+ triton_sem_t unlink_semaphore; //MUST be used EXCLUSIVELY for unlink
+ ssm_me match_entry;
+ union
+ {
+ unsigned char transport_address_data[MAX_RAW_ADDRESS_SIZE]; /*the actual data size depends on transport_type.
+ And this whole union is expected to be no bigger
+ than this array
+ */
+ ssmptcp_addrargs_t tcp_args;
+ //TODO: When they become available, add other ssm_transport address_arg types in this union.
+ //For now, only TCP is defined
+ };
+ int made_from_deserialization; //1 if made from deserialization
+ triton_list_t ssm_transaction_list; //transaction list
+ transaction_t* array_of_tx; //this holds the same transactions in ssm_transaction_list. It is kept to avoid
+ //myriads of malloc/free calls. See end_epoch
+}mem_pub_t;
+
+
+static int mem_reg_predicate(struct triton_list_link *ll, void *base_address)
+{
+ mem_reg_t* mem_reg = triton_list_get_entry(ll, struct mem_reg_t, list_link);
+ return mem_reg->base_address == base_address;
+}
+
+static mem_reg_t* find_mem_reg(void* base_address)
+{
+ mem_reg_t *mem_reg = NULL;
+ triton_list_link_t* ll = NULL;
+ if(base_address)
+ ll = triton_list_find(&ssm_global_data.mem_reg_list, mem_reg_predicate, base_address);
+ if(ll)
+ mem_reg = triton_list_get_entry(ll, struct mem_reg_t, list_link);
+ return mem_reg;
+}
+
+static triton_ret_t register_memory_internal(void* base,
+ uint64_t size,
+ mem_reg_t** handle
+ )
+{
+ mem_reg_t *mem_reg = find_mem_reg(base);
+ if(mem_reg)
+ {
+ if(mem_reg->is_registered)
+ return TRITON_ERR_MEM_REGISTERED;
+ else if(mem_reg->is_internally_allocated)
+ {
+ /*The buffer was allocated through the net-module and its handle can just be returned*/
+
+ mem_reg->is_registered = TRUE;
+ mem_reg->transport_type = default_net_endpoint->transport_type;
+ *handle = mem_reg;
+ return TRITON_SUCCESS;
+ }
+ else
+ triton_assert(0) //this should never happen
+ }
+
+ mem_reg = malloc(sizeof(mem_reg_t));
+ if(!mem_reg)
+ return TRITON_ERR_NOMEM;
+ mem_reg->size = size;
+ mem_reg->publish_count = 0;
+ if(base)
+ {
+ mem_reg->base_address = base;
+ mem_reg->is_registered = TRUE;
+ mem_reg->is_internally_allocated = FALSE;
+ mem_reg->transport_type = default_net_endpoint->transport_type;
+ }
+ else //this call is made from buffer_allocate
+ {
+ base = malloc(size);
+ if(!base)
+ {
+ triton_list_del(&mem_reg->list_link);
+ free(mem_reg);
+ return TRITON_ERR_NOMEM;
+ }
+ mem_reg->base_address = base;
+ mem_reg->is_registered = FALSE;
+ mem_reg->is_internally_allocated = TRUE;
+ }
+ triton_list_link_clear(&mem_reg->list_link);
+ triton_mutex_lock(&ssm_global_data.mutex);
+ triton_list_add_back(&mem_reg->list_link, &ssm_global_data.mem_reg_list);
+ triton_mutex_unlock(&ssm_global_data.mutex);
+ *handle = mem_reg;
+ return TRITON_SUCCESS;
+}
+
+
+static triton_ret_t unregister_memory_internal(mem_reg_t* mem_reg_info)
+{
+ triton_list_link_t *ptr, *scratch;
+ mem_reg_t *mem_reg = NULL;
+
+ triton_list_for_each(ptr, scratch, &ssm_global_data.mem_reg_list)
+ {
+ mem_reg = triton_list_get_entry(ptr, struct mem_reg_t, list_link);
+ if(mem_reg == mem_reg_info)
+ break;
+ }
+ if(!mem_reg)
+ return TRITON_ERR_INVAL;
+ if(mem_reg->publish_count != 0)
+ return TRITON_ERR_MEM_IN_USE;
+
+ if(mem_reg->is_internally_allocated)
+ {
+ if(mem_reg->is_registered) //This is really unregister
+ {
+ mem_reg->is_registered = FALSE;
+ return TRITON_SUCCESS;
+ }
+ else //this is buffer_free
+ {
+ free(mem_reg->base_address);
+ }
+ }
+ triton_mutex_lock(&ssm_global_data.mutex);
+ triton_list_del(&mem_reg->list_link);
+ triton_mutex_unlock(&ssm_global_data.mutex);
+ free(mem_reg);
+ return TRITON_SUCCESS;
+}
+
+static void* buffer_allocate(size_t size)
+{
+ mem_reg_t* handle = NULL;
+ triton_ret_t ret = register_memory_internal(NULL, (uint64_t)size, &handle);
+ if(ret != TRITON_SUCCESS)
+ return NULL;
+ return handle->base_address;
+}
+
+static void buffer_free(void* buffer)
+{
+ triton_ret_t tret;
+ mem_reg_t *mem_reg = find_mem_reg(buffer);
+ if(!mem_reg)
+ {
+ PRINT_ERR_MSG("The buffer was not allocated through the net-module");
+ return;
+ }
+
+ if(mem_reg->is_registered)
+ {
+ PRINT_ERR_MSG("Trying to free a registered memory");
+ return;
+ }
+ tret = unregister_memory_internal(mem_reg);
+ triton_assert(tret);
+}
+
+/**
+ * Register memory for the communication subsystem.
+ */
+static triton_ret_t register_memory(void* base, //the base address of the memory region
+ size_t size, // the size in bytes of the memory region
+ triton_mem_reg_handle_t* handle /*handle returned by
+ the memory registration*/
+ )
+{
+ if(size == 0)
+ {
+ *handle = ZERO_SIZE_MEM_HANDLE;
+ return TRITON_SUCCESS;
+ }
+ return register_memory_internal(base, (uint64_t)size, (mem_reg_t**)(handle));
+}
+
+/**
+ * Unregister memory.
+ */
+static triton_ret_t unregister_memory(triton_mem_reg_handle_t handle /*handle previously returned
+ by the memory registration*/
+ )
+{
+ if(handle == ZERO_SIZE_MEM_HANDLE)
+ return TRITON_SUCCESS;
+ return unregister_memory_internal((mem_reg_t*)(((void*)(handle))));
+}
+
+
+static void target_side_wait_callback(void* callback_args, void* event_data)
+{
+ ssm_result result = (ssm_result)event_data;
+ mem_pub_t* mem_pub = (mem_pub_t*)callback_args;
+ if(result->op == SSM_OP_PUT || result->op == SSM_OP_GET)
+ triton_sem_up(&mem_pub->comm_semaphore);
+ else if(result->op == SSM_OP_UNLINK)
+ triton_sem_up(&mem_pub->unlink_semaphore);
+
+}
+
+
+/**
+ * Expose a registered memory portion to remote peers.
+ */
+static __blocking triton_ret_t publish_memory(triton_mem_reg_handle_t in_handle, //handle of a registered memory
+ size_t offset, /*offset (in bytes) from the base address
+ where the memory publishing should start from*/
+ size_t size, //size (in bytes) of the published area
+ triton_mem_pub_handle_t* out_handle /*handle returned from
+ publishing the memory*/
+ )
+{
+ int ret;
+ if(size == 0)
+ {
+ *out_handle = ZERO_SIZE_MEM_HANDLE;
+ return TRITON_SUCCESS;
+ }
+
+ mem_reg_t* mem_reg_info = (mem_reg_t*)in_handle;
+ mem_pub_t* mem_pub = NULL;
+
+ if(!mem_reg_info->is_registered)
+ return TRITON_ERR_INVAL;
+
+ if(mem_reg_info->size < (uint64_t)offset ||
+ (uint64_t)size > mem_reg_info->size - (uint64_t)offset)
+ return TRITON_ERR_MEM_OUT_OF_RANGE;
+
+ mem_pub = (mem_pub_t*)calloc(1, sizeof(mem_pub_t));
+ if(!mem_pub)
+ return TRITON_ERR_NOMEM;
+
+ mem_pub->base = (void*)(offset + (uint64_t)mem_reg_info->base_address);
+ mem_pub->size = size;
+ mem_pub->mem_reg_info = mem_reg_info;
+ mem_pub->transport_type = mem_reg_info->transport_type;
+ mem_pub->transport_instance = default_net_endpoint->ssm_instance;
+
+ triton_mutex_lock(&ssm_global_data.mutex);
+ mem_reg_info->publish_count++;
+ mem_pub->match_bits = ssm_global_data.one_sided_match_bits_counter++;
+ if(ssm_global_data.one_sided_match_bits_counter == MAX_ONE_SIDED_MATCH_BITS)
+ ssm_global_data.one_sided_match_bits_counter = 0;
+ triton_mutex_unlock(&ssm_global_data.mutex);
+
+ mem_pub->callback.pcb = target_side_wait_callback;
+ mem_pub->callback.cbdata = mem_pub;
+
+ mem_pub->mr = ssm_mr_create(NULL, mem_pub->base, mem_pub->size);
+ mem_pub->match_entry = ssm_link(default_net_endpoint->ssm_instance, mem_pub->match_bits, 0,
+ SSM_POS_TAIL, NULL, &mem_pub->callback, SSM_NOF);
+ ret = ssm_post(default_net_endpoint->ssm_instance, mem_pub->match_entry, mem_pub->mr, SSM_NOF);
+
+ if(ret)
+ {
+ ssm_unlink(mem_pub->transport_instance, mem_pub->match_entry);
+ free(mem_pub);
+ return TRITON_ERR_SSM_UNKNOWN;
+ }
+
+ triton_sem_init(&mem_pub->comm_semaphore, 0);
+ triton_sem_init(&mem_pub->unlink_semaphore, 0);
+
+ *out_handle = (uint64_t)mem_pub;
+
+ return TRITON_SUCCESS;
+}
+
+/**
+ * End access permissions to a memory portion previously exposed.
+ */
+static __blocking triton_ret_t unpublish_memory(triton_mem_pub_handle_t handle /*handle previously
+ returned by a memory publishing*/
+ )
+{
+ if(handle == ZERO_SIZE_MEM_HANDLE)
+ return TRITON_SUCCESS;
+
+ mem_pub_t* mem_pub = (mem_pub_t*)handle;
+ mem_reg_t* mem_reg_info = mem_pub->mem_reg_info;
+
+ if(mem_pub->made_from_deserialization)
+ return TRITON_ERR_INVAL; //this handle is not allowed to be released through unpublish
+
+ ssm_unlink(default_net_endpoint->ssm_instance, mem_pub->match_entry);
+ triton_sem_count_down_no_cancel(&mem_pub->unlink_semaphore, 1);
+ ssm_mr_destroy(mem_pub->mr);
+ triton_sem_destroy(&mem_pub->comm_semaphore);
+ triton_sem_destroy(&mem_pub->unlink_semaphore);
+ free(mem_pub);
+
+ triton_mutex_lock(&ssm_global_data.mutex);
+ mem_reg_info->publish_count--;
+ triton_mutex_unlock(&ssm_global_data.mutex);
+
+ return TRITON_SUCCESS;
+}
+
+
+static int get_serialized_handle_size(triton_serialized_handle_type_t handle_type, void* handle)
+{
+ switch(handle_type)
+ {
+ case MEM_PUB_HANDLE_TYPE:
+ return (int)(sizeof(((mem_pub_t*)0)->size) + sizeof(((mem_pub_t*)0)->match_bits)
+ + default_net_endpoint->transport_abstraction->get_addr_serialization_size());
+ default:
+ triton_assert("unknown triton_serialized_handle_type_t");
+ }
+ return -1; //this should never happen
+}
+
+
+static triton_ret_t get_serialized_handle(triton_serialized_handle_type_t handle_type,
+ void* pointer_to_in_handle,
+ char* serialized_handle)
+{
+/*
+ void get_serialized_address(ssm_Itp ssm_transport, char* address_data, char* serialization,
+ char* field_types_in_serialization, int* nb_fields_in_serialization);
+ */
+
+ int offset = 0;
+
+ mem_pub_t* mem_pub = NULL;
+ switch(handle_type)
+ {
+ case MEM_PUB_HANDLE_TYPE:
+ mem_pub = *((mem_pub_t**)pointer_to_in_handle);
+ serialized_handle[0] = (char)mem_pub->mem_reg_info->transport_type;
+ offset+=sizeof(char);
+ *((uint64_t*)(serialized_handle+offset)) = aehton64(mem_pub->size);
+ offset+=sizeof(uint64_t);
+ *((int*)(serialized_handle+offset)) = aehton64(mem_pub->match_bits);
+ offset+=sizeof(int);
+ default_net_endpoint->transport_abstraction->get_serialized_address(
+ default_net_endpoint->ssm_transport,
+ default_net_endpoint->raw_address, (uint8_t*)(serialized_handle+offset));
+ break;
+ default:
+ return TRITON_ERR_INVAL;
+ }
+ return TRITON_SUCCESS;
+}
+
+
+static void origin_side_wait_callback(void* callback_args, void* event_data)
+{
+ ssm_result result = (ssm_result)event_data;
+ mem_pub_t* mem_pub = (mem_pub_t*)callback_args;
+ transaction_t *pos, *scratch;
+ if((result->status == SSM_ST_COMPLETE && (result->op == SSM_OP_PUT || result->op == SSM_OP_GET)) ||
+ result->status == SSM_ST_CANCEL)
+ {
+ triton_mutex_lock(&mem_pub->mutex);
+ triton_list_for_each_entry(pos, scratch, &mem_pub->ssm_transaction_list, transaction_t, list_link)
+ {
+ if(pos->transaction == result->tx)
+ {
+ triton_list_del(&pos->list_link);
+ ssm_md_release(pos->md);
+ ssm_mr_destroy(pos->mr);
+ triton_sem_up(&mem_pub->comm_semaphore);
+ break;
+ }
+ }
+ triton_mutex_unlock(&mem_pub->mutex);
+ }
+}
+
+
+static triton_ret_t create_handle_from_serialization(triton_serialized_handle_type_t handle_type,
+ void* pointer_to_out_handle,
+ char* serialized_handle)
+{
+ mem_pub_t* mem_pub = NULL;
+ ssm_transport_type_t transport_type = serialized_handle[0];
+ ssm_transport_abstraction_t* transport_abstraction = get_transport_abstraction(transport_type);
+ triton_assert(transport_abstraction);
+ if(!transport_abstraction)
+ return TRITON_ERR_SSM_TRSPT_TYPE_UNAVAILABLE;
+
+ serialized_handle += sizeof(char);
+ switch(handle_type)
+ {
+ case MEM_PUB_HANDLE_TYPE:
+ mem_pub = (mem_pub_t*)calloc(1, sizeof(mem_pub_t));
+ if(!mem_pub)
+ return TRITON_ERR_NOMEM;
+ mem_pub->made_from_deserialization = 1;
+ serialized_handle += sizeof(char);
+ mem_pub->size = aentoh64(*((uint64_t*)serialized_handle));
+ serialized_handle += sizeof(uint64_t);
+ mem_pub->match_bits = aentoh32(*((ssm_bits*)serialized_handle));
+ serialized_handle += sizeof(ssm_bits);
+
+ mem_pub->callback.pcb = origin_side_wait_callback;
+ mem_pub->callback.cbdata = mem_pub;
+ mem_pub->array_of_tx = NULL;
+
+ triton_list_init(&mem_pub->ssm_transaction_list);
+
+ triton_mutex_init(&mem_pub->mutex, NULL);
+ triton_sem_init(&mem_pub->comm_semaphore, 0);
+
+ switch(transport_type)
+ {
+ case SSM_TCP:
+ transport_abstraction->create_address_arg_data((uint8_t*)serialized_handle, mem_pub->transport_address_data);
+ break;
+ default:
+ ;
+ }
+
+ *((triton_mem_pub_handle_t*)pointer_to_out_handle) = (triton_mem_pub_handle_t)mem_pub;
+ break;
+ default:
+ return TRITON_ERR_INVAL;
+ }
+ return TRITON_SUCCESS;
+}
+
+static triton_ret_t free_serialized_handle(triton_serialized_handle_type_t handle_type,
+ void* pointer_to_in_handle)
+{
+ mem_pub_t* mem_pub = NULL;
+ switch(handle_type)
+ {
+ case MEM_PUB_HANDLE_TYPE:
+ mem_pub = *((mem_pub_t**)pointer_to_in_handle);
+
+ if(!mem_pub->made_from_deserialization)
+ return TRITON_ERR_INVAL; //it is forbidden to call this function on a handle
+ //that was not created through serialization
+ if(mem_pub->array_of_tx)
+ free(mem_pub->array_of_tx);
+
+ triton_mutex_destroy(&mem_pub->mutex);
+ triton_sem_destroy(&mem_pub->comm_semaphore);
+ free(mem_pub);
+ *((triton_mem_pub_handle_t*)pointer_to_in_handle) = 0;
+ break;
+ default:
+ return TRITON_ERR_INVAL;
+ }
+ return TRITON_SUCCESS;
+}
+
+
+static void* progress_thread_func(void* args)
+{
+ network_endpoint_t* net_endpoint;
+ struct timeval tv;
+ tv.tv_sec = 0;
+ tv.tv_usec = SSM_WAIT_TIMEOUT;
+ while(!ssm_global_data.exiting)
+ {
+ if(ssm_global_data.net_endpoint_list.count == 1) //only 1 ssm_instance
+ ssm_wait(default_net_endpoint->ssm_instance, NULL);
+ else //several instances
+ {
+ triton_list_link_t *ptr, *scratch;
+ triton_list_for_each(ptr, scratch, &ssm_global_data.net_endpoint_list)
+ {
+ net_endpoint = triton_list_get_entry(ptr, network_endpoint_t, list_link);
+ ssm_wait(net_endpoint->ssm_instance, &tv);
+ }
+ }
+ }
+ return NULL;
+}
+
+typedef enum one_sided_op_t
+{
+ TRITON_SSM_OP_GET,
+ TRITON_SSM_OP_PUT
+}one_sided_op_t;
+
+static triton_ret_t issue_one_sided_op(
+ triton_mem_pub_handle_t handle,
+ triton_iov_item_t* iovs,
+ char** dest_bufs,
+ int iov_items_count,
+ one_sided_op_t one_sided_op
+ )
+{
+ int i;
+ mem_pub_t* mem_pub = (mem_pub_t*)handle;
+ transaction_t* transactions = NULL;
+ transactions = malloc(iov_items_count*sizeof(transaction_t));
+ if(!transactions)
+ return TRITON_ERR_NOMEM;
+
+ for(i=0; i<iov_items_count; i++)
+ {
+ transactions[i].md = ssm_md_add(NULL, iovs[i].offset, iovs[i].size);
+ transactions[i].mr = ssm_mr_create(NULL, dest_bufs[i], iovs[i].size);
+ if(one_sided_op == TRITON_SSM_OP_GET)
+ {
+ transactions[i].transaction = ssm_get(mem_pub->transport_instance, mem_pub->initial_owner_address,
+ transactions[i].md, transactions[i].mr, mem_pub->match_bits, &mem_pub->callback, SSM_NOF);
+ }
+ else
+ {
+ transactions[i].transaction = ssm_put(mem_pub->transport_instance, mem_pub->initial_owner_address,
+ transactions[i].mr, transactions[i].md, mem_pub->match_bits, &mem_pub->callback, SSM_NOF);
+ }
+ triton_list_link_clear(&transactions[i].list_link);
+ triton_mutex_lock(&mem_pub->mutex);
+ triton_list_add_back(&transactions[i].list_link, &mem_pub->ssm_transaction_list);
+ triton_mutex_unlock(&mem_pub->mutex);
+ mem_pub->array_of_tx = transactions;
+ }
+ return TRITON_SUCCESS;
+}
+
+static __blocking triton_ret_t get(triton_msg_method_addr_t from,
+ triton_mem_pub_handle_t handle,
+ triton_iov_item_t* iovs,
+ char** dest_bufs,
+ int iov_items_count)
+{
+ return issue_one_sided_op(handle, iovs, dest_bufs, iov_items_count, TRITON_SSM_OP_GET);
+}
+
+static __blocking triton_ret_t put(triton_msg_method_addr_t to,
+ triton_mem_pub_handle_t handle,
+ triton_iov_item_t* iovs,
+ char** src_bufs,
+ int iov_items_count)
+{
+ return issue_one_sided_op(handle, iovs, src_bufs, iov_items_count, TRITON_SSM_OP_PUT);
+}
+
+
+static __blocking triton_ret_t start_epoch(triton_mem_pub_handle_t handle)
+{
+ mem_pub_t *mem_pub = (mem_pub_t*)handle;
+ network_endpoint_t *network_endpoint = get_network_endpoint(mem_pub->transport_type);
+ if(!network_endpoint)
+ return TRITON_ERR_SSM_TRSPT_TYPE_UNAVAILABLE;
+ mem_pub->transport_instance = network_endpoint->ssm_instance;
+ mem_pub->initial_owner_address = network_endpoint->transport_abstraction->create_addr_handle(
+ network_endpoint->address_interface, (void*)mem_pub->transport_address_data);
+ return TRITON_SUCCESS;
+}
+
+
+static __blocking triton_ret_t end_epoch(triton_mem_pub_handle_t handle)
+{
+ int i;
+ triton_ret_t tret;
+ int count;
+ mem_pub_t *mem_pub = (mem_pub_t*)handle;
+ transaction_t *pos, *scratch;
+ network_endpoint_t *network_endpoint = get_network_endpoint(mem_pub->transport_type);
+ triton_assert(network_endpoint);
+ if(!network_endpoint)
+ {
+ return TRITON_ERR_SSM_TRSPT_TYPE_UNAVAILABLE; //this should never happen.
+ }
+ tret = triton_sem_count_down(&mem_pub->comm_semaphore, mem_pub->ssm_transaction_list.count);
+ if(tret == TRITON_ERR_CANCELED)
+ {
+ triton_mutex_lock(&mem_pub->mutex);
+ count = mem_pub->ssm_transaction_list.count;
+ triton_list_for_each_entry(pos, scratch, &mem_pub->ssm_transaction_list, transaction_t, list_link)
+ {
+ ssm_cancel(mem_pub->transport_instance, pos->transaction);
+ }
+ triton_mutex_unlock(&mem_pub->mutex);
+ triton_sem_count_down_no_cancel(&mem_pub->comm_semaphore, count);
+ }
+ network_endpoint->address_interface->destroy(network_endpoint->address_interface,
+ mem_pub->initial_owner_address);
+ return tret;
+}
+
+
+static __blocking triton_ret_t wait(triton_mem_pub_handle_t mem_handle,
+ void* wait_info)
+{
+ triton_ret_t tret;
+ mem_pub_t* mem_pub = (mem_pub_t*)mem_handle;
+ mem_pub->nb_expected_completions = *((int*)wait_info);
+ return triton_sem_count_down(&mem_pub->comm_semaphore, mem_pub->nb_expected_completions);
+}
+
+static struct triton_msg_method triton_msg_ssm_method =
+{
+ .name = "ssm",
+
+ .addr_self = ssm_addr_self,
+ .addr_equal = ssm_addr_equal,
+ .addr_to_string = ssm_addr_to_string,
+ .addr_lookup = ssm_addr_lookup,
+ .addr_free = ssm_addr_free,
+ .addr_copy = ssm_addr_copy,
+ .new_tag = ssm_new_tag,
+
+ .group_add = ssm_group_add,
+ .send = ssm_send,
+ .recv = ssm_recv,
+ .recv_any = ssm_recv_any,
+
+ .set_one_sided_threshold = set_one_sided_threshold,
+ .get_one_sided_threshold = get_one_sided_threshold,
+ .buffer_allocate = buffer_allocate,
+ .buffer_free = buffer_free,
+ .register_memory = register_memory,
+ .unregister_memory = unregister_memory,
+ .publish_memory = publish_memory,
+ .unpublish_memory = unpublish_memory,
+ .get_serialized_handle_size = get_serialized_handle_size,
+ .create_handle_from_serialization = create_handle_from_serialization,
+ .get_serialized_handle = get_serialized_handle,
+ .free_serialized_handle = free_serialized_handle,
+ .get = get,
+ .put = put,
+ .start_epoch = start_epoch,
+ .end_epoch = end_epoch,
+ .wait = wait
+};
+
+
+triton_ret_t triton_msg_ssm_init(void)
+{
+
+ triton_ret_t tret;
+
+
+ triton_mutex_init(&ssm_global_data.mutex, NULL);
+ triton_mutex_init(&ssm_global_data.two_sided_mutex, NULL);
+ triton_list_init(&ssm_global_data.mem_reg_list);
+ triton_list_init(&ssm_global_data.net_endpoint_list);
+ triton_list_init(&ssm_global_data.generic_recv_data_list);
+
+ triton_sem_init_thread_safe(&ssm_global_data.mutex);
+
+ init_ethernet();
+ /*Init the other available tyransports here (e.g. init_ib()*/
+
+ set_default_net_endpoint();
+
+ assert(group_table == NULL);
+ group_table = triton_hash_init(ssm_group_compare, triton_hash_32bit_hash, 1024);
+ if(!group_table)
+ {
+ triton_err(triton_log_default, "%s:%d: Failed to initialize group hashtable", __FILE__, __LINE__);
+ assert(group_table);
+ }
+
+ pthread_create(&ssm_global_data.progress_thread, NULL, progress_thread_func, NULL);
+ tret = triton_msg_method_register(&triton_msg_ssm_method);
+ return tret;
+}
+
+void triton_msg_ssm_finalize(void)
+{
+ ssm_global_data.exiting = TRUE;
+ triton_sem_finalize_thread_safe(&ssm_global_data.mutex);
+ triton_msg_method_unregister("ssm");
+
+ release_ethernet();
+ /*Release the other available transports here (e.g. release_ib()*/
+
+ if(group_table != NULL)
+ {
+ triton_hash_destroy_and_finalize(group_table, struct triton_ssm_group, link, group_entry_free);
+ group_table = NULL;
+ }
+
+ pthread_join(ssm_global_data.progress_thread, NULL);
+ triton_mutex_destroy(&ssm_global_data.mutex);
+ triton_mutex_destroy(&ssm_global_data.two_sided_mutex);
+}
+
+__attribute__((constructor)) void triton_msg_ssm_init_register(void);
+
+__attribute__((constructor)) void triton_msg_ssm_init_register(void)
+{
+ triton_init_register("triton.net.ssm", triton_msg_ssm_init, triton_msg_ssm_finalize, NULL, "triton.net.message", "triton.resource.ssm");
+}
+
+
+//TEST
+struct triton_msg_method* get_msg_method()
+{
+ return &triton_msg_ssm_method;
+}
diff --git a/code/src/net/ssm/ssm-method.h b/code/src/net/ssm/ssm-method.h
new file mode 100644
index 0000000..031e4e5
--- /dev/null
+++ b/code/src/net/ssm/ssm-method.h
@@ -0,0 +1,8 @@
+#ifndef __SSM_METHOD_H__
+#define __SSM_METHOD_H__
+#include "src/common/triton-error.h"
+
+triton_ret_t triton_msg_ssm_init(void);
+void triton_msg_ssm_finalize(void);
+
+#endif //__SSM_METHOD_H__
diff --git a/code/src/net/ssm/ssm-tcp-transport.c b/code/src/net/ssm/ssm-tcp-transport.c
new file mode 100644
index 0000000..ae0b8a3
--- /dev/null
+++ b/code/src/net/ssm/ssm-tcp-transport.c
@@ -0,0 +1,161 @@
+#include "ssm-transport.h"
+#include "src/remote/byteswap.h"
+#include <string.h>
+
+#define JZ_DEBUG
+static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once)
+{
+ static int count = 0;
+
+ if(count && once)
+ return;
+ count++;
+ printf("process ( pid = %d | rank = %d ) is waiting in %s at %s:%d for debugger\n",
+ getpid(), rank, function, file, line);
+ fflush(stdout);
+ for(;;);
+}
+
+#ifdef JZ_DEBUG
+ void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once);
+ #define JZ_WAIT_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0)
+ #define JZ_WAIT_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0);}while(0)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1);}while(0)
+#else
+ #define JZ_WAIT_FOR_DEBUGGER()
+ #define JZ_WAIT_FOR_DEBUGGER_IF()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF()
+#endif
+static uint64_t get_address(uint8_t* raw_address_data, void* additional_data)
+{
+ /*
+ The address is made of an IP address and a port
+ */
+ uint64_t addr = ((uint64_t)raw_address_data[0]) << (40) | //IP
+ ((uint64_t)raw_address_data[1]) << (32)| //IP
+ ((uint64_t)raw_address_data[2]) << (24)| //IP
+ ((uint64_t)raw_address_data[3]) << (16); //IP
+
+#if __BYTE_ORDER == __LITTLE_ENDIAN
+ addr |= (((uint64_t)raw_address_data[5]) << (8) | //port
+ ((uint64_t)raw_address_data[4])); //port
+#else
+ addr |= (((uint64_t)raw_address_data[4]) << (8) | //port
+ ((uint64_t)raw_address_data[5])); //port
+#endif
+
+ /*NOTE: The port is expected to be in the right node endianess in raw_address_data*/
+
+ (void)additional_data;
+ return addr;
+}
+
+unsigned char *g_raw_data;
+static uint64_t get_address_from_str_addr(char* str_addr, void* additional_data)
+{
+ uint8_t raw_address[6];
+ (void)additional_data;
+ sscanf(str_addr, "tcp::%hhu.%hhu.%hhu.%hhu|%hu", &raw_address[0], &raw_address[1], &raw_address[2],
+ &raw_address[3], ((unsigned short*)(&raw_address[4])));
+ g_raw_data = raw_address;
+ return get_address(raw_address, NULL);
+}
+
+static void get_str_address(uint64_t address, char str_address[], void* additional_args)
+{
+ uint8_t ip[4];
+ unsigned short port = address&0xffff;
+ address >>= 16;
+ ip[3]=(uint8_t)(address&0xff);
+ address >>= 8;
+ ip[2]=(uint8_t)(address&0xff);
+ address >>= 8;
+ ip[1]=(uint8_t)(address&0xff);
+ address >>= 8;
+ ip[0]=(uint8_t)(address&0xff);
+
+ sprintf(str_address, "tcp::%d.%d.%d.%d|%d",(uint8_t)ip[0], (uint8_t)ip[1], (uint8_t)ip[2],
+ (uint8_t)ip[3], port);
+ (void)additional_args;
+}
+
+
+static inline int get_address_data_size()
+{
+ return 6; //6 bytes ... assuming IPv4
+}
+
+static ssm_Itp create_transport(void* args)
+{
+ return NULL;
+}
+
+static inline int get_addr_serialization_size()
+{
+ return sizeof(uint8_t) + sizeof(uint8_t)*4 + sizeof(short); //assuming IPv4
+}
+
+static void get_serialized_address(ssm_Itp transport, uint8_t* address_data, uint8_t* serialization)
+{
+ //IP
+ serialization[0] = address_data[0];
+ serialization[1] = address_data[1];
+ serialization[2] = address_data[2];
+ serialization[3] = address_data[3];
+
+ //port
+ *(int16_t*)(&(serialization[4])) = aehton16(*((short*)(&address_data[4])));
+}
+
+
+static void get_raw_address_in_host_bytes(ssm_Itp transport, uint8_t* net_bytes_address, uint8_t* host_bytes)
+{
+ //IP
+ host_bytes[0] = net_bytes_address[0];
+ host_bytes[1] = net_bytes_address[1];
+ host_bytes[2] = net_bytes_address[2];
+ host_bytes[3] = net_bytes_address[3];
+
+ //port
+ *(int16_t*)(&(host_bytes[4])) = aentoh16(*((short*)(&net_bytes_address[4])));
+}
+
+static void create_address_arg_data(uint8_t* raw_address_data, void* transport_address_arg_out)
+{
+ unsigned short port;
+ ssmptcp_addrargs_t* tcp_args = (ssmptcp_addrargs_t*)transport_address_arg_out;
+ tcp_args->host = (uint8_t*)transport_address_arg_out+sizeof(ssmptcp_addrargs_t);
+ sprintf(tcp_args->host, "%d.%d.%d.%d", (int)raw_address_data[0], (int)raw_address_data[1],
+ (int)raw_address_data[2], (int)raw_address_data[3]);
+ port = *((unsigned short*)(&raw_address_data[4]));
+ tcp_args->port = aentoh16(port);
+}
+
+static ssm_Haddr create_addr_handle(ssm_Iaddr addr_interface, void* addr_args)
+{
+ ssmptcp_addrargs_t* tcp_args = (ssmptcp_addrargs_t*)addr_args;
+ return addr_interface->create(addr_interface, tcp_args);
+}
+
+static ssm_transport_abstraction_t tcp_abstraction =
+{
+ .get_address = get_address,
+ .get_address_from_str_addr = get_address_from_str_addr,
+ .get_address_data_size = get_address_data_size,
+ .get_str_address = get_str_address,
+ .create_transport = create_transport,
+ .get_addr_serialization_size = get_addr_serialization_size,
+ .get_serialized_address = get_serialized_address,
+ .get_raw_address_in_host_bytes = get_raw_address_in_host_bytes,
+ .create_address_arg_data = create_address_arg_data,
+ .create_addr_handle = create_addr_handle
+};
+
+ssm_transport_abstraction_t* get_tcp_abstraction()
+{
+ return &tcp_abstraction;
+}
diff --git a/code/src/net/ssm/ssm-transport.h b/code/src/net/ssm/ssm-transport.h
new file mode 100644
index 0000000..fa48166
--- /dev/null
+++ b/code/src/net/ssm/ssm-transport.h
@@ -0,0 +1,69 @@
+#ifndef __SSM_TRANSPORT_H__
+#define __SSM_TRANSPORT_H__
+
+/*
+This file defines a transport specific interface for SSM use in triton
+*/
+
+#include "ssm.h"
+#include "ssmptcp.h"
+
+#define SSM_MAX_ADDR_ARG_SIZE 256
+
+typedef enum ssm_transport_type_t
+{
+ SSM_TCP,
+ SSM_UDP,
+ SSM_IB,
+ /*
+ etc.
+ */
+}ssm_transport_type_t;
+
+typedef struct ssm_transport_abstraction_t
+{
+ /*
+ get address from raw address data component laid out as a uint8_t[]
+ */
+ uint64_t (*get_address)(uint8_t* address_data, void* additional_data);
+
+ uint64_t (*get_address_from_str_addr)(char* str_addr, void* additional_data);
+
+ /*
+ get the number of items in the uint8_t address_data[] representation of the address
+ */
+ int (*get_address_data_size)();
+
+ /*
+ get the string representation of the address; i.e. the ssm://stuff_stuff form of the address
+ additional_args contains anything required by the transport to generate the
+ str_address from address
+ */
+ void (*get_str_address)(uint64_t address, char* str_address, void* additional_args);
+
+ ssm_Itp (*create_transport)(void*); //the prototype of this function will change for SURE. For now it will take a void*
+
+ int (*get_addr_serialization_size)();
+
+ /*Just as the name says.
+ The serialization is in network byte order (See the TCP implementation for an example)
+ */
+ void (*get_serialized_address)(ssm_Itp transport, uint8_t* address_data, uint8_t* serialization);
+
+ /*Just as the name says.
+ The serialization is in network byte order (See the TCP implementation for an example)
+ */
+ void (*get_raw_address_in_host_bytes)(ssm_Itp transport, uint8_t* net_bytes_address,
+ uint8_t* host_bytes);
+
+ /*transport_address_arg_out must be a buffer of at least SSM_MAX_ADDR_ARG_SIZE
+ */
+ void (*create_address_arg_data)(uint8_t* raw_address_data, void* transport_address_arg_out);
+ ssm_Haddr (*create_addr_handle)(ssm_Iaddr addr_interface, void* addr_args);
+
+}ssm_transport_abstraction_t;
+
+
+ssm_transport_abstraction_t* get_tcp_abstraction();
+
+#endif //__SSM_TRANSPORT_H__
diff --git a/code/src/net/tests/module.mk.in b/code/src/net/tests/module.mk.in
index 0316fbf..07e68b4 100644
--- a/code/src/net/tests/module.mk.in
+++ b/code/src/net/tests/module.mk.in
@@ -4,7 +4,9 @@ DIR := src/net/tests
ifneq (,$(BUILD_MPI))
AETESTSRC += $(DIR)/mpi-client-server.ae\
- $(DIR)/mpi-one-sided.ae
+ $(DIR)/mpi-one-sided.ae\
+ $(DIR)/ssm-one-sided.ae\
+ $(DIR)/ssm-sr-test.ae
MODCFLAGS_$(DIR)/mpi-client-server = $(MPICFLAGS)
MODLDFLAGS_$(DIR)/mpi-client-server = $(MPILDFLAGS)
diff --git a/code/src/net/tests/mpi-one-sided.ae b/code/src/net/tests/mpi-one-sided.ae
index 8b5a4ae..4b25e6d 100644
--- a/code/src/net/tests/mpi-one-sided.ae
+++ b/code/src/net/tests/mpi-one-sided.ae
@@ -30,11 +30,13 @@ do
#define MB 1048576
#define PRINT(...) do{printf(__VA_ARGS__); fflush(stdout);}while(0)
+#define TRACE(...) do{PRINT("%s:%d[%s]", __FILE__, __LINE__, __FUNCTION__); PRINT(__VA_ARGS__);}while(0)
+#define TRACE_IF(cond,...) do{if((cond)) {TRACE(__VA_ARGS__);}}while(0)
#define MAX_NB_SERVERS 256 /* For I'm using byte offsets and writing/reading char to check the test results.
I'm writing ranks to know who wrote where and who read from where; so
- more than 255 is troublesome. Writing char is just for the readability
- of the results of this test.
+ more than 255 is troublesome to visually assess the correctness of the test.
+ Writing char is just for the readability of the results of this test.
*/
@@ -187,9 +189,9 @@ static void init_server_data()
triton_ret_t tret;
if(!server_data)
{
- server_data = triton_buffer_allocate(group, MAX_DATA_SIZE);
+ server_data = triton_buffer_allocate(self, MAX_DATA_SIZE);
triton_assert(server_data);
- tret = triton_register_memory(group, server_data, MAX_DATA_SIZE, &server_mrh);
+ tret = triton_register_memory(self, server_data, MAX_DATA_SIZE, &server_mrh);
triton_assert(tret == TRITON_SUCCESS);
}
memset(server_data, rank, MAX_DATA_SIZE);
@@ -202,9 +204,9 @@ static void free_server_data()
triton_ret_t tret;
if(server_data)
{
- tret = triton_unregister_memory(group, server_mrh);
+ tret = triton_unregister_memory(self, server_mrh);
triton_assert(tret == TRITON_SUCCESS);
- triton_buffer_free(group, server_data);
+ triton_buffer_free(self, server_data);
server_data = NULL;
}
}
@@ -299,6 +301,7 @@ static void get_io_address_extent(triton_iov_item_t* iovs, int nb_iov_items,
out_mem_pub_params->size = largest_end - smallest_start;
}
+#if 0
int *g_rpc_size;
char *g_rpc_packet;
triton_iov_item_t* g_iovs;
@@ -307,6 +310,8 @@ char* g_serialized_mem_pub;
char* g_str_client_addr;
uint64_t g_unique_token;
uint32_t *g_offset;
+int g_serialized_mem_pub_size;
+#endif
/*
unique_token is used to identify a specific request. Two pairs (str_client_addr, unique_token)
@@ -326,15 +331,15 @@ static int build_io_rpc(triton_iov_item_t* iovs, int* valid_iov_indices, int nb_
int rpc_size = 0;
- g_rpc_size = &rpc_size;
- g_rpc_packet = rpc_packet;
-
strcpy(rpc_packet, str_client_addr);
rpc_size += (strlen(str_client_addr)+1);
memcpy(rpc_packet + rpc_size, &unique_token, sizeof(uint64_t));
rpc_size+=sizeof(uint64_t);
+ memcpy(rpc_packet + rpc_size, &serialized_mem_pub_size, sizeof(int));
+ rpc_size+=sizeof(int);
+
memcpy(rpc_packet + rpc_size, serialized_mem_pub, serialized_mem_pub_size);
rpc_size += serialized_mem_pub_size;
@@ -349,21 +354,26 @@ static int build_io_rpc(triton_iov_item_t* iovs, int* valid_iov_indices, int nb_
static void extract_io_rpc_data(char* rpc_packet, triton_iov_item_t** iovs, int *nb_iov_items,
- char** serialized_mem_pub, char **str_client_addr, uint64_t *unique_token)
+ int *serialized_mem_pub_size, char** serialized_mem_pub, char **str_client_addr, uint64_t *unique_token)
{
uint32_t offset = 0;
- g_offset = &offset;
-
- uint32_t mem_pub_h_size = (uint32_t)triton_get_serialized_handle_size(group, MEM_PUB_HANDLE_TYPE);
- *str_client_addr= rpc_packet;
+ *str_client_addr = rpc_packet;
offset += (strlen(*str_client_addr) + 1);
+
*unique_token = *(uint64_t*)(rpc_packet + offset);
offset += sizeof(uint64_t);
- *serialized_mem_pub = rpc_packet + offset;
- offset += mem_pub_h_size;
+
+
+ *serialized_mem_pub_size = *((int*)(rpc_packet + offset));
+ offset += sizeof(int);
+
+ *serialized_mem_pub = (char*)(rpc_packet + offset);
+ offset += *serialized_mem_pub_size;
+
*nb_iov_items = *((int*)(rpc_packet + offset));
offset += sizeof(int);
+
*iovs = (triton_iov_item_t*)(rpc_packet + offset);
}
@@ -384,7 +394,6 @@ static __blocking void do_client_io(char* mem_start, triton_mem_reg_handle_t mrh
//the same data as the equivalent RPC payload
char* packet_ptr = NULL; //just to make gcc happy; see where it is used
int serialized_mem_pub_h_size;
- triton_wait_handle_t wait_handle;
uint32_t rpc_size = 0;
uint32_t byte_received = 0;
@@ -394,11 +403,11 @@ static __blocking void do_client_io(char* mem_start, triton_mem_reg_handle_t mrh
for(i=0; i<nb_iov_items; i++)
valid_iov_indices[i] = i;
- tret = triton_publish_memory(group, mrh, mem_pub_handle_params.offset, mem_pub_handle_params.size, &mph);
+ tret = triton_publish_memory(self, mrh, mem_pub_handle_params.offset, mem_pub_handle_params.size, &mph);
triton_assert(tret == TRITON_SUCCESS);
- serialized_mem_pub_h_size = triton_get_serialized_handle_size(group, MEM_PUB_HANDLE_TYPE);
- tret = triton_get_serialized_handle(group, MEM_PUB_HANDLE_TYPE, (void*)&mph, some_buf);
+ serialized_mem_pub_h_size = triton_get_serialized_handle_size(self, MEM_PUB_HANDLE_TYPE, &mph);
+ tret = triton_get_serialized_handle(self, MEM_PUB_HANDLE_TYPE, (void*)&mph, some_buf);
triton_assert(tret == TRITON_SUCCESS);
rpc_size = build_io_rpc(iovs, valid_iov_indices, nb_iov_items, some_buf, serialized_mem_pub_h_size,
@@ -419,17 +428,10 @@ static __blocking void do_client_io(char* mem_start, triton_mem_reg_handle_t mrh
triton_assert(tret == TRITON_SUCCESS);
triton_assert(rpc_packet[0] == TEST_TYPE_ACKS);
- tret = triton_set_serialized_handle(group, WAIT_HANDLE_TYPE, (void*)(&wait_handle),
- rpc_packet);
- triton_assert(tret == TRITON_SUCCESS);
-
- tret = triton_wait(group, mph, wait_handle);
- triton_assert(tret == TRITON_SUCCESS);
-
- tret = triton_free_serialized_handle(group, WAIT_HANDLE_TYPE, (void*)&wait_handle);
+ tret = triton_wait(self, mph, (void*)(rpc_packet+sizeof(char)));
triton_assert(tret == TRITON_SUCCESS);
- tret = triton_unpublish_memory(group, mph);
+ tret = triton_unpublish_memory(self, mph);
triton_assert(tret == TRITON_SUCCESS);
if(iot == IOT_READ)
@@ -472,10 +474,13 @@ to the begining of its iovs
static int has_iov()
{
char buf[1025];
- int addr_length = strlen(str_self);
+ char str_rank[5];
+ int addr_length;
+ sprintf(str_rank, "%d", rank);
+ addr_length = strlen(str_rank);
while(fgets(buf, 1024, input_file))
{
- if((strlen(buf) > addr_length + 2) && buf[0] == '<' && strncmp(buf+1, str_self, addr_length) == 0)
+ if((strlen(buf) > addr_length + 2) && buf[0] == '<' && strncmp(buf+1, str_rank, addr_length) == 0)
return 1;
}
return 0;
@@ -524,7 +529,7 @@ static int extract_iov(triton_iov_item_t** iov, io_type_t* iot)
input_extract_phase_t input_extract_phase = UNDEFINED;
int nb_extracted_iov_items = 0;
- sprintf(end_of_iov, "</%s>", str_self);
+ sprintf(end_of_iov, "</%d>", rank);
*iot = IOT_UNDEFINED;
while(fgets(line, 1024, input_file))
{
@@ -586,7 +591,7 @@ static __blocking void enter_client_loop()
memset(buffer, rank, MAX_DATA_SIZE); //initialize the client memory
- tret = triton_register_memory(group, buffer, MB, &mrh1);
+ tret = triton_register_memory(self, buffer, MB, &mrh1);
triton_assert(tret == TRITON_SUCCESS);
pwait
@@ -612,7 +617,7 @@ static __blocking void enter_client_loop()
}
}
- tret = triton_unregister_memory(group, mrh1);
+ tret = triton_unregister_memory(self, mrh1);
triton_assert(tret == TRITON_SUCCESS);
notify_for_client_shutdown();
@@ -647,11 +652,11 @@ static __blocking void do_server_io(triton_addr_t client_addr, triton_mem_pub_ha
if(test_type == TEST_TYPE_READ)
{
- tret = triton_msg_put(group, client_addr, mem_pub_handle,iov, bufs, nb_iov_items);
+ tret = triton_msg_put(self, client_addr, mem_pub_handle,iov, bufs, nb_iov_items);
}
else
{
- tret = triton_msg_get(group, client_addr, mem_pub_handle, iov, bufs, nb_iov_items);
+ tret = triton_msg_get(self, client_addr, mem_pub_handle, iov, bufs, nb_iov_items);
}
triton_assert(tret == TRITON_SUCCESS);
free(bufs);
@@ -715,7 +720,6 @@ static __blocking void enter_server_loop()
volatile int request_is_received;
triton_mutex_t mutex;
triton_mutex_t rpc_mutex;
- int serialized_mem_pub_h_size = triton_get_serialized_handle_size(group, MEM_PUB_HANDLE_TYPE);
char listener_rpc_packet[MAX_RPC_SIZE];
char* listener_packet_ptr = listener_rpc_packet;
@@ -724,6 +728,7 @@ static __blocking void enter_server_loop()
pwait
{
+ pprivate int serialized_mem_pub_h_size;
pprivate char rpc_packet[MAX_RPC_SIZE];
char* packet_ptr = rpc_packet; //just to make gcc happy
pprivate triton_iov_item_t* iovs;
@@ -763,7 +768,7 @@ static __blocking void enter_server_loop()
}
/*
- The server wait for the last spawned pbranch to copy the rpc data before it goes
+ The server waits for the last spawned pbranch to copy the rpc data before it goes
on listening to the next one again. The wait is necessary since the same
buffer is used.
*/
@@ -794,7 +799,7 @@ static __blocking void enter_server_loop()
{
case TEST_TYPE_READ:
case TEST_TYPE_WRITE:
- extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub,
+ extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub_h_size, &serialized_mem_pub,
&str_client_addr, &unique_token);
PRINT("Primary Server %s handling I/O %lu\n", str_self, unique_token);
client_addr = private_requester_addr;
@@ -830,7 +835,7 @@ static __blocking void enter_server_loop()
}
/*This server sends its share of responses now*/
- tret = triton_set_serialized_handle(group, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
+ tret = triton_create_handle_from_serialization(self, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
serialized_mem_pub);
triton_assert(tret == TRITON_SUCCESS);
nb_iov_items_for_specific_server = get_response_map_for_specific_server(
@@ -842,11 +847,11 @@ static __blocking void enter_server_loop()
}
PRINT("Primary Server %s about to serve I/O %lu\n",
str_self, unique_token);
- tret = triton_start_epoch(group, mem_pub_handle);
+ tret = triton_start_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
do_server_io(client_addr, mem_pub_handle, partial_iov,
nb_iov_items_for_specific_server, rpc_packet[0] /*READ or WRITE*/);
- tret = triton_end_epoch(group, mem_pub_handle);
+ tret = triton_end_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
PRINT("Primary Server %s done serving I/O %lu\n",
str_self, unique_token);
@@ -855,7 +860,7 @@ static __blocking void enter_server_loop()
display_server_side_read_data(server_data, partial_iov,
nb_iov_items_for_specific_server, str_client_addr);
}
- tret = triton_free_serialized_handle(group, MEM_PUB_HANDLE_TYPE,
+ tret = triton_free_serialized_handle(self, MEM_PUB_HANDLE_TYPE,
(void*)(&mem_pub_handle));
triton_assert(tret == TRITON_SUCCESS);
if(nb_scheduled_servers == 1)
@@ -878,28 +883,28 @@ static __blocking void enter_server_loop()
//Handle pipeline response requests from some primarily contacted server
case TEST_TYPE_READ_SCHEDULE_TRANSFER_REQUEST:
case TEST_TYPE_WRITE_SCHEDULE_TRANSFER_REQUEST:
- extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub,
+ extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub_h_size, &serialized_mem_pub,
&str_client_addr, &unique_token);
PRINT("Auxiliary Server %s received pipelined request for I/O %lu\n",
str_self, unique_token);
- tret = triton_set_serialized_handle(group, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
+ tret = triton_create_handle_from_serialization(self, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
serialized_mem_pub);
triton_assert(tret == TRITON_SUCCESS);
tret = triton_addr_lookup(str_client_addr, &client_addr);
triton_assert(tret);
- tret = triton_start_epoch(group, mem_pub_handle);
+ tret = triton_start_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
do_server_io(client_addr, mem_pub_handle, iovs, nb_iov_items,
rpc_packet[0] == TEST_TYPE_READ_SCHEDULE_TRANSFER_REQUEST?
TEST_TYPE_READ:TEST_TYPE_WRITE);
- tret = triton_end_epoch(group, mem_pub_handle);
+ tret = triton_end_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
if(rpc_packet[0] == TEST_TYPE_WRITE_SCHEDULE_TRANSFER_REQUEST)
{
display_server_side_read_data(server_data, iovs, nb_iov_items, str_client_addr);
}
- tret = triton_free_serialized_handle(group, MEM_PUB_HANDLE_TYPE,
+ tret = triton_free_serialized_handle(self, MEM_PUB_HANDLE_TYPE,
(void*)(&mem_pub_handle));
triton_assert(tret == TRITON_SUCCESS);
rpc_packet[0]=TEST_TYPE_ACK_FRAGMENT;
diff --git a/code/src/net/tests/sample_one_sided_test_input b/code/src/net/tests/sample_one_sided_test_input
index ef594d2..a2eef1e 100644
--- a/code/src/net/tests/sample_one_sided_test_input
+++ b/code/src/net/tests/sample_one_sided_test_input
@@ -11,7 +11,7 @@
# Empty lines are allowed; they are just ignored
#
-<mpi://0>
+<0>
<READ>
#The bnumber right below is the number of iov items in this read
8
@@ -34,9 +34,9 @@
200,1
1000,10
</WRITE>
-</mpi://0>
+<0>
-<mpi://4>
+<4>
<READ>
2
2048,1024
@@ -50,4 +50,4 @@
850,1
900,1
</WRITE>
-</mpi://4>
+<4>
diff --git a/code/src/net/tests/mpi-one-sided.ae b/code/src/net/tests/ssm-one-sided.ae
similarity index 83%
copy from code/src/net/tests/mpi-one-sided.ae
copy to code/src/net/tests/ssm-one-sided.ae
index 8b5a4ae..c073e0a 100644
--- a/code/src/net/tests/mpi-one-sided.ae
+++ b/code/src/net/tests/ssm-one-sided.ae
@@ -1,9 +1,10 @@
+#include "src/net/ssm/ssm-method.h"
+
#include "src/aesop/aesop.h"
#include <mpi.h>
#include <stdio.h>
#include <errno.h>
#include "../triton-message-method.hae"
-#include "../mpi/mpi-method.h"
#include <time.h>
#include <stdlib.h>
#include <string.h>
@@ -11,7 +12,41 @@
#include <ctype.h>
#include <unistd.h>
#include "src/aesop/aesop-support.hae"
+#include <stdlib.h>
+
+#define JZ_DEBUG
+static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once)
+{
+ static int count = 0;
+
+ if(count && once)
+ return;
+ count++;
+ printf("process ( pid = %d | rank = %d ) is waiting in %s at %s:%d for debugger\n",
+ getpid(), rank, function, file, line);
+ fflush(stdout);
+ for(;;);
+}
+
+#ifdef JZ_DEBUG
+ void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once);
+ #define JZ_WAIT_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0)
+ #define JZ_WAIT_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0);}while(0)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1);}while(0)
+#else
+ #define JZ_WAIT_FOR_DEBUGGER()
+ #define JZ_WAIT_FOR_DEBUGGER_IF()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF()
+#endif
+
+#define ENV_ETHERNET_IP "ETHERNET_IP"
+#define ENV_ETHERNET_PORT "ETHERNET_PORT"
+#define ETHERNET_LOOPBACK "127.0.0.1"
#define MIN(a,b) ((a)<=(b)?(a):(b))
#define MAX(a,b) ((a)>=(b)?(a):(b))
@@ -30,6 +65,8 @@ do
#define MB 1048576
#define PRINT(...) do{printf(__VA_ARGS__); fflush(stdout);}while(0)
+#define TRACE(...) do{PRINT("%s:%d[%s]", __FILE__, __LINE__, __FUNCTION__); PRINT(__VA_ARGS__);}while(0)
+#define TRACE_IF(cond,...) do{if((cond)) {TRACE(__VA_ARGS__);}}while(0)
#define MAX_NB_SERVERS 256 /* For I'm using byte offsets and writing/reading char to check the test results.
I'm writing ranks to know who wrote where and who read from where; so
@@ -38,7 +75,7 @@ do
*/
-
+#define PORT_BASE 5000
#define MAX_DATA_SIZE 50*MB
int rank = -1;
@@ -121,11 +158,25 @@ static inline int is_local_process_a_server()
return is_server(rank);
}
+static void set_environment()
+{
+ int ret;
+ char env_string[101];
+ sprintf(env_string, "%s=%s", ENV_ETHERNET_IP, ETHERNET_LOOPBACK);
+ ret = putenv(env_string);
+ triton_assert(ret == 0);
+ sprintf(env_string, "%s=%d", ENV_ETHERNET_PORT, PORT_BASE+rank);
+ ret = putenv(env_string);
+ triton_assert(ret == 0);
+}
+
static __blocking void open_endpoints()
{
- self = triton_addr_self("mpi");
- sprintf(str_self, "mpi://%d", rank);
- triton_msg_group_add("mpi", "test-one-sided", &group);
+ triton_string_t tstring;
+ self = triton_addr_self("ssm");
+ triton_addr_to_string(self, &tstring);
+ strcpy(str_self, tstring.string);
+ triton_msg_group_add("ssm", "test-one-sided", &group);
tag = triton_msg_new_tag(group);
}
@@ -152,7 +203,7 @@ static __blocking void discover_servers(void)
{
if(is_server(i))
{
- sprintf(str_addr, "mpi://%d", i);
+ sprintf(str_addr, "ssm://tcp::%s|%d:", ETHERNET_LOOPBACK, PORT_BASE+i);
tret = triton_addr_lookup(str_addr, &server_addrs[nb_servers++]);
triton_assert(tret == TRITON_SUCCESS);
}
@@ -187,9 +238,9 @@ static void init_server_data()
triton_ret_t tret;
if(!server_data)
{
- server_data = triton_buffer_allocate(group, MAX_DATA_SIZE);
+ server_data = triton_buffer_allocate(self, MAX_DATA_SIZE);
triton_assert(server_data);
- tret = triton_register_memory(group, server_data, MAX_DATA_SIZE, &server_mrh);
+ tret = triton_register_memory(self, server_data, MAX_DATA_SIZE, &server_mrh);
triton_assert(tret == TRITON_SUCCESS);
}
memset(server_data, rank, MAX_DATA_SIZE);
@@ -202,9 +253,9 @@ static void free_server_data()
triton_ret_t tret;
if(server_data)
{
- tret = triton_unregister_memory(group, server_mrh);
+ tret = triton_unregister_memory(self, server_mrh);
triton_assert(tret == TRITON_SUCCESS);
- triton_buffer_free(group, server_data);
+ triton_buffer_free(self, server_data);
server_data = NULL;
}
}
@@ -274,7 +325,11 @@ static __blocking void notify_for_client_shutdown()
uint32_t packet_size = sizeof(char);
for(i=0; i<nb_servers; i++)
{
- tret = triton_msg_send(server_addrs[i], group, tag, 1, &payload, &packet_size);
+ do
+ {
+ tret = triton_msg_send(server_addrs[i], group, tag, 1, &payload, &packet_size);
+ }
+ while(tret != TRITON_SUCCESS);
triton_assert(tret == TRITON_SUCCESS);
}
}
@@ -299,6 +354,7 @@ static void get_io_address_extent(triton_iov_item_t* iovs, int nb_iov_items,
out_mem_pub_params->size = largest_end - smallest_start;
}
+#if 0
int *g_rpc_size;
char *g_rpc_packet;
triton_iov_item_t* g_iovs;
@@ -307,6 +363,8 @@ char* g_serialized_mem_pub;
char* g_str_client_addr;
uint64_t g_unique_token;
uint32_t *g_offset;
+int g_serialized_mem_pub_size;
+#endif
/*
unique_token is used to identify a specific request. Two pairs (str_client_addr, unique_token)
@@ -326,15 +384,15 @@ static int build_io_rpc(triton_iov_item_t* iovs, int* valid_iov_indices, int nb_
int rpc_size = 0;
- g_rpc_size = &rpc_size;
- g_rpc_packet = rpc_packet;
-
strcpy(rpc_packet, str_client_addr);
rpc_size += (strlen(str_client_addr)+1);
memcpy(rpc_packet + rpc_size, &unique_token, sizeof(uint64_t));
rpc_size+=sizeof(uint64_t);
+ memcpy(rpc_packet + rpc_size, &serialized_mem_pub_size, sizeof(int));
+ rpc_size+=sizeof(int);
+
memcpy(rpc_packet + rpc_size, serialized_mem_pub, serialized_mem_pub_size);
rpc_size += serialized_mem_pub_size;
@@ -349,21 +407,26 @@ static int build_io_rpc(triton_iov_item_t* iovs, int* valid_iov_indices, int nb_
static void extract_io_rpc_data(char* rpc_packet, triton_iov_item_t** iovs, int *nb_iov_items,
- char** serialized_mem_pub, char **str_client_addr, uint64_t *unique_token)
+ int *serialized_mem_pub_size, char** serialized_mem_pub, char **str_client_addr, uint64_t *unique_token)
{
uint32_t offset = 0;
- g_offset = &offset;
-
- uint32_t mem_pub_h_size = (uint32_t)triton_get_serialized_handle_size(group, MEM_PUB_HANDLE_TYPE);
- *str_client_addr= rpc_packet;
+ *str_client_addr = rpc_packet;
offset += (strlen(*str_client_addr) + 1);
+
*unique_token = *(uint64_t*)(rpc_packet + offset);
offset += sizeof(uint64_t);
- *serialized_mem_pub = rpc_packet + offset;
- offset += mem_pub_h_size;
+
+
+ *serialized_mem_pub_size = *((int*)(rpc_packet + offset));
+ offset += sizeof(int);
+
+ *serialized_mem_pub = (char*)(rpc_packet + offset);
+ offset += *serialized_mem_pub_size;
+
*nb_iov_items = *((int*)(rpc_packet + offset));
offset += sizeof(int);
+
*iovs = (triton_iov_item_t*)(rpc_packet + offset);
}
@@ -384,7 +447,6 @@ static __blocking void do_client_io(char* mem_start, triton_mem_reg_handle_t mrh
//the same data as the equivalent RPC payload
char* packet_ptr = NULL; //just to make gcc happy; see where it is used
int serialized_mem_pub_h_size;
- triton_wait_handle_t wait_handle;
uint32_t rpc_size = 0;
uint32_t byte_received = 0;
@@ -394,11 +456,11 @@ static __blocking void do_client_io(char* mem_start, triton_mem_reg_handle_t mrh
for(i=0; i<nb_iov_items; i++)
valid_iov_indices[i] = i;
- tret = triton_publish_memory(group, mrh, mem_pub_handle_params.offset, mem_pub_handle_params.size, &mph);
+ tret = triton_publish_memory(self, mrh, mem_pub_handle_params.offset, mem_pub_handle_params.size, &mph);
triton_assert(tret == TRITON_SUCCESS);
- serialized_mem_pub_h_size = triton_get_serialized_handle_size(group, MEM_PUB_HANDLE_TYPE);
- tret = triton_get_serialized_handle(group, MEM_PUB_HANDLE_TYPE, (void*)&mph, some_buf);
+ serialized_mem_pub_h_size = triton_get_serialized_handle_size(self, MEM_PUB_HANDLE_TYPE, &mph);
+ tret = triton_get_serialized_handle(self, MEM_PUB_HANDLE_TYPE, (void*)&mph, some_buf);
triton_assert(tret == TRITON_SUCCESS);
rpc_size = build_io_rpc(iovs, valid_iov_indices, nb_iov_items, some_buf, serialized_mem_pub_h_size,
@@ -411,7 +473,12 @@ static __blocking void do_client_io(char* mem_start, triton_mem_reg_handle_t mrh
packet_ptr = rpc_packet;
PRINT("[CLIENT %s] about to issue %s request %lu\n", str_self, iot == IOT_READ ?"READ":"WRITE",
(uint64_t)mph);
- tret = triton_msg_send(primary_server_addr, group, tag, 1, &packet_ptr, &rpc_size);
+
+ do
+ {
+ tret = triton_msg_send(primary_server_addr, group, tag, 1, &packet_ptr, &rpc_size);
+ }
+ while(tret != TRITON_SUCCESS); //retries
triton_assert(tret == TRITON_SUCCESS);
//receive ack
@@ -419,17 +486,10 @@ static __blocking void do_client_io(char* mem_start, triton_mem_reg_handle_t mrh
triton_assert(tret == TRITON_SUCCESS);
triton_assert(rpc_packet[0] == TEST_TYPE_ACKS);
- tret = triton_set_serialized_handle(group, WAIT_HANDLE_TYPE, (void*)(&wait_handle),
- rpc_packet);
- triton_assert(tret == TRITON_SUCCESS);
-
- tret = triton_wait(group, mph, wait_handle);
- triton_assert(tret == TRITON_SUCCESS);
-
- tret = triton_free_serialized_handle(group, WAIT_HANDLE_TYPE, (void*)&wait_handle);
+ tret = triton_wait(self, mph, (void*)(rpc_packet+sizeof(char)));
triton_assert(tret == TRITON_SUCCESS);
- tret = triton_unpublish_memory(group, mph);
+ tret = triton_unpublish_memory(self, mph);
triton_assert(tret == TRITON_SUCCESS);
if(iot == IOT_READ)
@@ -464,6 +524,7 @@ static void close_input()
fclose(input_file);
}
+
/*
Returns 1 if this process has iov. 0 otherwise.
If this process has iovs, then the file read pointer is set
@@ -472,10 +533,13 @@ to the begining of its iovs
static int has_iov()
{
char buf[1025];
- int addr_length = strlen(str_self);
+ char str_rank[5];
+ int addr_length;
+ sprintf(str_rank, "%d", rank);
+ addr_length = strlen(str_rank);
while(fgets(buf, 1024, input_file))
{
- if((strlen(buf) > addr_length + 2) && buf[0] == '<' && strncmp(buf+1, str_self, addr_length) == 0)
+ if((strlen(buf) > addr_length + 2) && buf[0] == '<' && strncmp(buf+1, str_rank, addr_length) == 0)
return 1;
}
return 0;
@@ -524,7 +588,7 @@ static int extract_iov(triton_iov_item_t** iov, io_type_t* iot)
input_extract_phase_t input_extract_phase = UNDEFINED;
int nb_extracted_iov_items = 0;
- sprintf(end_of_iov, "</%s>", str_self);
+ sprintf(end_of_iov, "</%d>", rank);
*iot = IOT_UNDEFINED;
while(fgets(line, 1024, input_file))
{
@@ -568,7 +632,6 @@ static int extract_iov(triton_iov_item_t** iov, io_type_t* iot)
return nb_extracted_iov_items;
}
-
/*
*/
static __blocking void enter_client_loop()
@@ -586,7 +649,7 @@ static __blocking void enter_client_loop()
memset(buffer, rank, MAX_DATA_SIZE); //initialize the client memory
- tret = triton_register_memory(group, buffer, MB, &mrh1);
+ tret = triton_register_memory(self, buffer, MB, &mrh1);
triton_assert(tret == TRITON_SUCCESS);
pwait
@@ -612,7 +675,7 @@ static __blocking void enter_client_loop()
}
}
- tret = triton_unregister_memory(group, mrh1);
+ tret = triton_unregister_memory(self, mrh1);
triton_assert(tret == TRITON_SUCCESS);
notify_for_client_shutdown();
@@ -647,11 +710,11 @@ static __blocking void do_server_io(triton_addr_t client_addr, triton_mem_pub_ha
if(test_type == TEST_TYPE_READ)
{
- tret = triton_msg_put(group, client_addr, mem_pub_handle,iov, bufs, nb_iov_items);
+ tret = triton_msg_put(self, client_addr, mem_pub_handle,iov, bufs, nb_iov_items);
}
else
{
- tret = triton_msg_get(group, client_addr, mem_pub_handle, iov, bufs, nb_iov_items);
+ tret = triton_msg_get(self, client_addr, mem_pub_handle, iov, bufs, nb_iov_items);
}
triton_assert(tret == TRITON_SUCCESS);
free(bufs);
@@ -715,7 +778,6 @@ static __blocking void enter_server_loop()
volatile int request_is_received;
triton_mutex_t mutex;
triton_mutex_t rpc_mutex;
- int serialized_mem_pub_h_size = triton_get_serialized_handle_size(group, MEM_PUB_HANDLE_TYPE);
char listener_rpc_packet[MAX_RPC_SIZE];
char* listener_packet_ptr = listener_rpc_packet;
@@ -724,6 +786,7 @@ static __blocking void enter_server_loop()
pwait
{
+ pprivate int serialized_mem_pub_h_size;
pprivate char rpc_packet[MAX_RPC_SIZE];
char* packet_ptr = rpc_packet; //just to make gcc happy
pprivate triton_iov_item_t* iovs;
@@ -763,7 +826,7 @@ static __blocking void enter_server_loop()
}
/*
- The server wait for the last spawned pbranch to copy the rpc data before it goes
+ The server waits for the last spawned pbranch to copy the rpc data before it goes
on listening to the next one again. The wait is necessary since the same
buffer is used.
*/
@@ -794,7 +857,7 @@ static __blocking void enter_server_loop()
{
case TEST_TYPE_READ:
case TEST_TYPE_WRITE:
- extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub,
+ extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub_h_size, &serialized_mem_pub,
&str_client_addr, &unique_token);
PRINT("Primary Server %s handling I/O %lu\n", str_self, unique_token);
client_addr = private_requester_addr;
@@ -824,13 +887,17 @@ static __blocking void enter_server_loop()
TEST_TYPE_WRITE_SCHEDULE_TRANSFER_REQUEST;
rpc_size++;
io_packet_ptr = io_rpc_packets[j];
- tret = triton_msg_send(server_addrs[i], group, tag, 1, &io_packet_ptr, &rpc_size);
+ do
+ {
+ tret = triton_msg_send(server_addrs[i], group, tag, 1, &io_packet_ptr, &rpc_size);
+ }
+ while(tret != TRITON_SUCCESS);
triton_assert(tret == TRITON_SUCCESS);
j++;
}
/*This server sends its share of responses now*/
- tret = triton_set_serialized_handle(group, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
+ tret = triton_create_handle_from_serialization(self, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
serialized_mem_pub);
triton_assert(tret == TRITON_SUCCESS);
nb_iov_items_for_specific_server = get_response_map_for_specific_server(
@@ -842,11 +909,11 @@ static __blocking void enter_server_loop()
}
PRINT("Primary Server %s about to serve I/O %lu\n",
str_self, unique_token);
- tret = triton_start_epoch(group, mem_pub_handle);
+ tret = triton_start_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
do_server_io(client_addr, mem_pub_handle, partial_iov,
nb_iov_items_for_specific_server, rpc_packet[0] /*READ or WRITE*/);
- tret = triton_end_epoch(group, mem_pub_handle);
+ tret = triton_end_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
PRINT("Primary Server %s done serving I/O %lu\n",
str_self, unique_token);
@@ -855,7 +922,7 @@ static __blocking void enter_server_loop()
display_server_side_read_data(server_data, partial_iov,
nb_iov_items_for_specific_server, str_client_addr);
}
- tret = triton_free_serialized_handle(group, MEM_PUB_HANDLE_TYPE,
+ tret = triton_free_serialized_handle(self, MEM_PUB_HANDLE_TYPE,
(void*)(&mem_pub_handle));
triton_assert(tret == TRITON_SUCCESS);
if(nb_scheduled_servers == 1)
@@ -867,8 +934,12 @@ static __blocking void enter_server_loop()
PRINT("Primary Server %s about to acknowledge I/O %lu\n",
str_self, unique_token);
packet_ptr = rpc_packet;
- tret = triton_msg_send(client_addr, group, tag, 1,
- &packet_ptr, &rpc_size);
+ do
+ {
+ tret = triton_msg_send(client_addr, group, tag, 1,
+ &packet_ptr, &rpc_size);
+ }
+ while(tret != TRITON_SUCCESS);
triton_assert(tret == TRITON_SUCCESS);
PRINT("Primary Server %s completed and acknowledged I/O %lu\n",
str_self, unique_token);
@@ -878,36 +949,40 @@ static __blocking void enter_server_loop()
//Handle pipeline response requests from some primarily contacted server
case TEST_TYPE_READ_SCHEDULE_TRANSFER_REQUEST:
case TEST_TYPE_WRITE_SCHEDULE_TRANSFER_REQUEST:
- extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub,
+ extract_io_rpc_data(&rpc_packet[1], &iovs, &nb_iov_items, &serialized_mem_pub_h_size, &serialized_mem_pub,
&str_client_addr, &unique_token);
PRINT("Auxiliary Server %s received pipelined request for I/O %lu\n",
str_self, unique_token);
- tret = triton_set_serialized_handle(group, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
+ tret = triton_create_handle_from_serialization(self, MEM_PUB_HANDLE_TYPE, (void*)(&mem_pub_handle),
serialized_mem_pub);
triton_assert(tret == TRITON_SUCCESS);
tret = triton_addr_lookup(str_client_addr, &client_addr);
triton_assert(tret);
- tret = triton_start_epoch(group, mem_pub_handle);
+ tret = triton_start_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
do_server_io(client_addr, mem_pub_handle, iovs, nb_iov_items,
rpc_packet[0] == TEST_TYPE_READ_SCHEDULE_TRANSFER_REQUEST?
TEST_TYPE_READ:TEST_TYPE_WRITE);
- tret = triton_end_epoch(group, mem_pub_handle);
+ tret = triton_end_epoch(self, mem_pub_handle);
triton_assert(tret == TRITON_SUCCESS);
if(rpc_packet[0] == TEST_TYPE_WRITE_SCHEDULE_TRANSFER_REQUEST)
{
display_server_side_read_data(server_data, iovs, nb_iov_items, str_client_addr);
}
- tret = triton_free_serialized_handle(group, MEM_PUB_HANDLE_TYPE,
+ tret = triton_free_serialized_handle(self, MEM_PUB_HANDLE_TYPE,
(void*)(&mem_pub_handle));
triton_assert(tret == TRITON_SUCCESS);
rpc_packet[0]=TEST_TYPE_ACK_FRAGMENT;
memcpy(&rpc_packet[1], &unique_token, sizeof(unique_token));
rpc_size = 1 + sizeof(unique_token);
packet_ptr = rpc_packet;
- tret = triton_msg_send(private_requester_addr/*primary server in this case*/, group,
- tag, 1, &packet_ptr, &rpc_size);
+ do
+ {
+ tret = triton_msg_send(private_requester_addr/*primary server in this case*/, group,
+ tag, 1, &packet_ptr, &rpc_size);
+ }
+ while(tret != TRITON_SUCCESS);
triton_assert(tret == TRITON_SUCCESS);
triton_addr_free(client_addr);
break;
@@ -926,7 +1001,11 @@ static __blocking void enter_server_loop()
memcpy(&rpc_packet[1], &unique_token, sizeof(unique_token));
rpc_size = 1 + sizeof(unique_token);
packet_ptr = rpc_packet;
- tret = triton_msg_send(wait_info->client_addr, group, tag, 1, &packet_ptr, &rpc_size);
+ do
+ {
+ tret = triton_msg_send(wait_info->client_addr, group, tag, 1, &packet_ptr, &rpc_size);
+ }
+ while(tret != TRITON_SUCCESS);
triton_assert(tret == TRITON_SUCCESS);
PRINT("Primary Server %s acknowledging I/O %lu after pipelined transfers\n",
str_self, unique_token);
@@ -967,9 +1046,17 @@ static void print_test_description()
__blocking int entry_point(int argc, char** argv)
{
int dummy;
+ triton_ret_t tret;
srand(time(NULL));
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
MPI_Comm_size(MPI_COMM_WORLD, &size);
+
+ set_environment();
+
+
+ JZ_WAIT_FOR_DEBUGGER();
+ tret = triton_msg_ssm_init();
+ triton_assert(tret == TRITON_SUCCESS);
if(argc != 2)
{
if(rank == 0)
@@ -1014,6 +1101,7 @@ __blocking int entry_point(int argc, char** argv)
close_trace_in_file();
close_input();
close_endpoints();
+ triton_msg_ssm_finalize();
return 0;
}
diff --git a/code/src/net/tests/ssm-sr-test.ae b/code/src/net/tests/ssm-sr-test.ae
new file mode 100644
index 0000000..5a09386
--- /dev/null
+++ b/code/src/net/tests/ssm-sr-test.ae
@@ -0,0 +1,233 @@
+#include "src/net/ssm/ssm-method.h"
+
+#include "src/aesop/aesop.h"
+#include <mpi.h>
+#include <stdio.h>
+#include <errno.h>
+#include "../triton-message-method.hae"
+#include <time.h>
+#include <stdlib.h>
+#include <string.h>
+#include <stdarg.h>
+#include <ctype.h>
+#include <unistd.h>
+#include "src/aesop/aesop-support.hae"
+#include <stdlib.h>
+
+
+#define JZ_DEBUG
+static __attribute__((noinline)) void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once)
+{
+ static int count = 0;
+
+ if(count && once)
+ return;
+ count++;
+ printf("process ( pid = %d | rank = %d ) is waiting in %s at %s:%d for debugger\n",
+ getpid(), rank, function, file, line);
+ fflush(stdout);
+ for(;;);
+}
+
+#ifdef JZ_DEBUG
+ void __jz_wait_for_debugger(const char* file, const char* function, int line, int rank, int once);
+ #define JZ_WAIT_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0)
+ #define JZ_WAIT_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 0);}while(0)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER() __jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1)
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF(cond) \
+ do{if((cond))__jz_wait_for_debugger(__FILE__, __FUNCTION__, __LINE__, -1, 1);}while(0)
+#else
+ #define JZ_WAIT_FOR_DEBUGGER()
+ #define JZ_WAIT_FOR_DEBUGGER_IF()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER()
+ #define JZ_WAIT_ONCE_FOR_DEBUGGER_IF()
+#endif
+
+#define ENV_ETHERNET_IP "ETHERNET_IP"
+#define ENV_ETHERNET_PORT "ETHERNET_PORT"
+#define ETHERNET_LOOPBACK "127.0.0.1"
+
+#define MIN(a,b) ((a)<=(b)?(a):(b))
+#define MAX(a,b) ((a)>=(b)?(a):(b))
+#define ABORT_ON_MPI_FAILURE(ret,err_msg_format,...) \
+do \
+{ \
+ if((ret) != MPI_SUCCESS) \
+ { \
+ triton_err(triton_log_default, "%s:%d: [MPI Proc rank %d]" err_msg_format, \
+ __FILE__, __LINE__, rank, ##__VA_ARGS__); \
+ MPI_Abort(MPI_COMM_WORLD, ret); \
+ } \
+}while(0)
+
+#define KB 1024
+#define MB 1048576
+
+#define PRINT(...) do{printf(__VA_ARGS__); fflush(stdout);}while(0)
+#define TRACE PRINT
+#define TRACE_IF(cond,...) do{if((cond)) {PRINT("%s:%d[%s]", __FILE__, __LINE__, __FUNCTION__); PRINT(__VA_ARGS__);}}while(0)
+
+#define MAX_NB_SERVERS 256 /* For I'm using byte offsets and writing/reading char to check the test results.
+ I'm writing ranks to know who wrote where and who read from where; so
+ more than 255 is troublesome. Writing char is just for the readability
+ of the results of this test.
+ */
+
+
+#define PORT_BASE 5000
+
+#define MAX_DATA_SIZE 50*MB
+int rank = -1;
+int size = -1;
+int nb_servers = -1;
+
+triton_addr_t self;
+char str_self[51]; //to containg for instance "mpi://0" for rank 0
+triton_msg_group_t group;
+triton_msg_tag_t tag;
+triton_addr_t server_addrs[MAX_NB_SERVERS];
+
+
+static int is_prime(int n)
+{
+ int i, half = n/2;
+ if(n == 0)
+ return 0;
+ if(n == 1)
+ return 1;
+ for(i=2; i<=half; i++)
+ if( (n%i) == 0)
+ return 0;
+ return 1;
+}
+
+static inline int is_server(int process_rank)
+{
+ if(size <=2 )
+ return process_rank >=size -1;
+ return is_prime(process_rank);
+}
+
+
+/*Just some arbitrary logic to decide who is client and who is server
+*/
+static inline int is_local_process_a_server()
+{
+ return is_server(rank);
+}
+
+static void set_environment()
+{
+ int ret;
+ char env_string[101];
+ sprintf(env_string, "%s=%s", ENV_ETHERNET_IP, ETHERNET_LOOPBACK);
+ ret = putenv(env_string);
+ triton_assert(ret == 0);
+ sprintf(env_string, "%s=%d", ENV_ETHERNET_PORT, PORT_BASE+rank);
+ ret = putenv(env_string);
+ triton_assert(ret == 0);
+}
+
+static __blocking void open_endpoints()
+{
+ triton_string_t tstring;
+ self = triton_addr_self("ssm");
+ triton_addr_to_string(self, &tstring);
+ strcpy(str_self, tstring.string);
+ triton_msg_group_add("ssm", "test-one-sided", &group);
+ tag = triton_msg_new_tag(group);
+}
+
+static __blocking void close_endpoints()
+{
+ int i;
+ triton_addr_free(self);
+ for(i=0; i< nb_servers; i++)
+ triton_addr_free(server_addrs[i]);
+}
+
+/*
+Discovering servers could be implemented differently; e.g. having the
+processes read from a file who's who
+*/
+static __blocking void discover_servers(void)
+{
+ triton_ret_t tret;
+ int ret;
+ int i;
+ nb_servers = 0;
+ char str_addr[51];
+ for(i=0; i< size; i++)
+ {
+ if(is_server(i))
+ {
+ sprintf(str_addr, "ssm://tcp::%s|%d", ETHERNET_LOOPBACK, PORT_BASE+i);
+ tret = triton_addr_lookup(str_addr, &server_addrs[nb_servers++]);
+ triton_assert(tret == TRITON_SUCCESS);
+ }
+ }
+}
+
+
+char buffer[MB];
+#define DATA_SIZE 100
+
+
+__blocking int entry_point(int argc, char** argv)
+{
+ int dummy;
+ int i;
+ triton_addr_t requester_addr;
+ triton_msg_tag_t unused_tag;
+ uint32_t max_recv_size;
+ uint32_t received_bytes;
+ char* buf = buffer;
+
+ uint32_t data_size = DATA_SIZE;
+ triton_ret_t tret;
+ srand(time(NULL));
+ MPI_Comm_rank(MPI_COMM_WORLD, &rank);
+ MPI_Comm_size(MPI_COMM_WORLD, &size);
+
+ set_environment();
+
+
+ tret = triton_msg_ssm_init();
+ triton_assert(tret == TRITON_SUCCESS);
+ open_endpoints();
+ discover_servers();
+
+ for(i=0; i<DATA_SIZE; i++)
+ buffer[i] = rank;
+
+ if(rank == 0)
+ {
+ sleep(1);
+ tret = triton_msg_send(server_addrs[0], group, tag, 1, &buf, &data_size);
+ triton_assert(tret == TRITON_SUCCESS)
+ for(i=0; i<DATA_SIZE; i++)
+ buffer[i] = rank+100;
+ tret = triton_msg_send(server_addrs[0], group, tag, 1, &buf, &data_size);
+ triton_assert(tret != TRITON_SUCCESS)
+ }
+ else
+ {
+ tret = triton_msg_recv_any(group, &requester_addr, &unused_tag, 1,
+ &buf, &max_recv_size, &received_bytes);
+ triton_assert(tret == TRITON_SUCCESS);
+ for(i=0; i<DATA_SIZE; i++)
+ printf("buffer[%d] = %d\n", i, buffer[i]);
+ tret = triton_msg_recv(requester_addr, group, unused_tag, 1,
+ &buf, &max_recv_size, &received_bytes);
+ triton_assert(tret == TRITON_SUCCESS);
+ for(i=0; i<DATA_SIZE; i++)
+ printf("buffer[%d] = %d\n", i, buffer[i]);
+ }
+
+ close_endpoints();
+ triton_msg_ssm_finalize();
+ return 0;
+}
+
+aesop_main_set_with_init(NULL, "triton.client", entry_point);
diff --git a/code/src/net/triton-message-method.hae b/code/src/net/triton-message-method.hae
index 3766f38..d292d6a 100644
--- a/code/src/net/triton-message-method.hae
+++ b/code/src/net/triton-message-method.hae
@@ -138,8 +138,10 @@ struct triton_msg_method
/*returns the size of buffer required to get the serialized representation
of the info type passed as argument
+ *pointer_to_handle must be compatible with handle_type
*/
- int (*get_serialized_handle_size)(triton_serialized_handle_type_t handle_type);
+ int (*get_serialized_handle_size)(triton_serialized_handle_type_t handle_type,
+ void* pointer_to_handle);
/*Ask the net-module to provide the serialized net-module-specific info
@@ -180,7 +182,7 @@ struct triton_msg_method
free_serialized_handle(...,h1) and free_serialized_handle(..., h1) might be wrong unless something like
free_serialized_handle(...,h1); h1 = h0; free_serialized_handle(..., h1) happens.
*/
- triton_ret_t (*set_serialized_handle)(triton_serialized_handle_type_t handle_type,
+ triton_ret_t (*create_handle_from_serialization)(triton_serialized_handle_type_t handle_type,
void* pointer_to_out_handle,
char* serialized_handle);
@@ -219,7 +221,9 @@ struct triton_msg_method
__blocking triton_ret_t (*end_epoch)(triton_mem_pub_handle_t in_handle);
- __blocking triton_ret_t (*wait)(triton_mem_pub_handle_t handle, triton_wait_handle_t wait_handle);
+ /*Wait_info is netmodule-specific data that defines the wait completion condition
+ */
+ __blocking triton_ret_t (*wait)(triton_mem_pub_handle_t handle, void* wait_info);
};
/**
diff --git a/code/src/net/triton-message.ae b/code/src/net/triton-message.ae
index 12bfc0e..a597cb7 100644
--- a/code/src/net/triton-message.ae
+++ b/code/src/net/triton-message.ae
@@ -397,10 +397,12 @@ triton_ret_t triton_addr_lookup(const char *name, triton_addr_t *addr)
*Set the message size beyond which one-sided communication is used instead 2-sided for message transfers.
* The size is in bytes
*/
-triton_ret_t triton_set_one_sided_threshold(triton_msg_group_t group, size_t size)
+triton_ret_t triton_set_one_sided_threshold(triton_addr_t local_addr, size_t size)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->set_one_sided_threshold(size);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->set_one_sided_threshold(size);
}
@@ -408,10 +410,12 @@ triton_ret_t triton_set_one_sided_threshold(triton_msg_group_t group, size_t siz
*Get the message size beyond which one-sided communication is used instead 2-sided for message transfers.
* The size is in bytes
*/
-triton_ret_t triton_get_one_sided_threshold(triton_msg_group_t group, size_t *size)
+triton_ret_t triton_get_one_sided_threshold(triton_addr_t local_addr, size_t *size)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->get_one_sided_threshold(size);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->get_one_sided_threshold(size);
}
/**
@@ -419,60 +423,72 @@ triton_ret_t triton_get_one_sided_threshold(triton_msg_group_t group, size_t *si
* call to pin a buffer.
* TODO: Should we include information about the type of buffer to allocate?
*/
-void* triton_buffer_allocate(triton_msg_group_t group, size_t size)
+void* triton_buffer_allocate(triton_addr_t local_addr, size_t size)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->buffer_allocate(size);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->buffer_allocate(size);
}
/**
* Free a buffer allocated with buffer_allocate.
*/
-void triton_buffer_free(triton_msg_group_t group, void *buffer)
+void triton_buffer_free(triton_addr_t local_addr, void *buffer)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- g->method->buffer_free(buffer);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ tm->buffer_free(buffer);
}
/**
* Register memory for the communication subsystem.
*/
-triton_ret_t triton_register_memory(triton_msg_group_t group, void* base,
+triton_ret_t triton_register_memory(triton_addr_t local_addr, void* base,
size_t size,
triton_mem_reg_handle_t* handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->register_memory(base, size, handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->register_memory(base, size, handle);
}
/**
* Unregister memory.
*/
-triton_ret_t triton_unregister_memory(triton_msg_group_t group, triton_mem_reg_handle_t handle)
+triton_ret_t triton_unregister_memory(triton_addr_t local_addr, triton_mem_reg_handle_t handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->unregister_memory(handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->unregister_memory(handle);
}
/**
* Expose a registered memory portion to remote peers.
*/
-__blocking triton_ret_t triton_publish_memory(triton_msg_group_t group,
+__blocking triton_ret_t triton_publish_memory(triton_addr_t local_addr,
triton_mem_reg_handle_t in_handle,
size_t offset,
size_t size,
triton_mem_pub_handle_t* out_handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->publish_memory(in_handle, offset, size, out_handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->publish_memory(in_handle, offset, size, out_handle);
}
/**
* End access permissions to a memory portion previously exposed.
*/
-__blocking triton_ret_t triton_unpublish_memory(triton_msg_group_t group, triton_mem_pub_handle_t handle)
+__blocking triton_ret_t triton_unpublish_memory(triton_addr_t local_addr, triton_mem_pub_handle_t handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->unpublish_memory(handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->unpublish_memory(handle);
}
@@ -480,11 +496,13 @@ __blocking triton_ret_t triton_unpublish_memory(triton_msg_group_t group, triton
/*returns the size of buffer required to get the serialized representation
of the info type passed as argument
*/
-int triton_get_serialized_handle_size(triton_msg_group_t group,
- triton_serialized_handle_type_t handle_type)
+int triton_get_serialized_handle_size(triton_addr_t local_addr,
+ triton_serialized_handle_type_t handle_type, void* pointer_to_handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->get_serialized_handle_size(handle_type);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->get_serialized_handle_size(handle_type, pointer_to_handle);
}
@@ -498,22 +516,25 @@ int triton_get_serialized_handle_size(triton_msg_group_t group,
2. serialized_handle must be large enough. A previous call to
get_serialized_handle_size is recommended to determine the required buffer size.
- 3. This function might succeed even if set_serialized_handle (below) was not called previously
+ 3. This function might succeed even if triton_create_handle_from_serialization (below) was not called previously
to return pointer_to_in_handle. E.g. When publish_memory returns a handle h, any
subsequent call to get_serialized_handle with h and handle_type == MEM_PUB_HANDLE_TYPE
will succeed.
*/
-triton_ret_t triton_get_serialized_handle(triton_msg_group_t group,
+triton_ret_t triton_get_serialized_handle(triton_addr_t local_addr,
triton_serialized_handle_type_t handle_type,
void* pointer_to_in_handle,
char* serialized_handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->get_serialized_handle(handle_type, pointer_to_in_handle, serialized_handle);
+
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->get_serialized_handle(handle_type, pointer_to_in_handle, serialized_handle);
}
/*
- Set an serialized info and get a handle to it.
+ Create a handle from serialized info and get a handle to it.
COMMENTS:
1. See comment 1 of get_serialized_handle above
@@ -531,13 +552,15 @@ triton_ret_t triton_get_serialized_handle(triton_msg_group_t group,
free_serialized_handle(...,h1) and free_serialized_handle(..., h1) might be wrong unless something like
free_serialized_handle(...,h1); h1 = h0; free_serialized_handle(..., h1) happens.
*/
-triton_ret_t triton_set_serialized_handle(triton_msg_group_t group,
+triton_ret_t triton_create_handle_from_serialization(triton_addr_t local_addr,
triton_serialized_handle_type_t handle_type,
void* pointer_to_out_handle,
char* serialized_handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->set_serialized_handle(handle_type, pointer_to_out_handle, serialized_handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->create_handle_from_serialization(handle_type, pointer_to_out_handle, serialized_handle);
}
/*
@@ -549,24 +572,25 @@ triton_ret_t triton_set_serialized_handle(triton_msg_group_t group,
with a triton_mem_pub_handle_t that was directly created through publish_memory. For such a handle
unpublish_memory is the right way to "free" things.
*/
-triton_ret_t triton_free_serialized_handle(triton_msg_group_t group, triton_serialized_handle_type_t handle_type,
+triton_ret_t triton_free_serialized_handle(triton_addr_t local_addr, triton_serialized_handle_type_t handle_type,
void* pointer_to_in_handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->free_serialized_handle(handle_type, pointer_to_in_handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->free_serialized_handle(handle_type, pointer_to_in_handle);
}
/*dest_bufs is a an array of buffer local to the initator of the operation
iov_items_count is the number of iov_items and also the number of bufs
*/
-__blocking triton_ret_t triton_msg_get(triton_msg_group_t group,
+__blocking triton_ret_t triton_msg_get(triton_addr_t local_addr,
triton_addr_t from,
triton_mem_pub_handle_t handle,
triton_iov_item_t* iovs,
char **dest_bufs,
int iov_items_count)
{
- /*The group is pased here for uniformity*/
struct triton_msg_method *tm;
triton_msg_method_addr_t maddr;
addr_split(from, &tm, &maddr);
@@ -576,14 +600,13 @@ __blocking triton_ret_t triton_msg_get(triton_msg_group_t group,
/*src_bufs is a an array of buffer local to the initator of the operation
iov_items_count is the number of iov_items and also the number of bufs
*/
-__blocking triton_ret_t triton_msg_put(triton_msg_group_t group,
+__blocking triton_ret_t triton_msg_put(triton_addr_t local_addr,
triton_addr_t to,
triton_mem_pub_handle_t handle,
triton_iov_item_t* iovs,
char ** src_bufs,
int iov_items_count)
{
- /*The group is pased here for uniformity*/
struct triton_msg_method *tm;
triton_msg_method_addr_t maddr;
addr_split(to, &tm, &maddr);
@@ -591,23 +614,29 @@ __blocking triton_ret_t triton_msg_put(triton_msg_group_t group,
}
-__blocking triton_ret_t triton_start_epoch(triton_msg_group_t group, triton_mem_pub_handle_t in_handle)
+__blocking triton_ret_t triton_start_epoch(triton_addr_t local_addr, triton_mem_pub_handle_t in_handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->start_epoch(in_handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->start_epoch(in_handle);
}
-__blocking triton_ret_t triton_end_epoch(triton_msg_group_t group, triton_mem_pub_handle_t in_handle)
+__blocking triton_ret_t triton_end_epoch(triton_addr_t local_addr, triton_mem_pub_handle_t in_handle)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->end_epoch(in_handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->end_epoch(in_handle);
}
-__blocking triton_ret_t triton_wait(triton_msg_group_t group, triton_mem_pub_handle_t handle,
- triton_wait_handle_t wait_handle)
+__blocking triton_ret_t triton_wait(triton_addr_t local_addr, triton_mem_pub_handle_t handle,
+ void* wait_info)
{
- struct triton_msg_group *g = (struct triton_msg_group *)group;
- return g->method->wait(handle, wait_handle);
+ struct triton_msg_method *tm;
+ triton_msg_method_addr_t maddr;
+ addr_split(local_addr, &tm, &maddr);
+ return tm->wait(handle, wait_info);
}
/*
diff --git a/code/src/net/triton-message.hae b/code/src/net/triton-message.hae
index 281caea..e482fcb 100644
--- a/code/src/net/triton-message.hae
+++ b/code/src/net/triton-message.hae
@@ -104,8 +104,6 @@ typedef uint64_t triton_mem_reg_handle_t;
typedef uint64_t triton_mem_pub_handle_t;
-typedef uint64_t triton_wait_handle_t;
-
typedef struct triton_iov_item_t
{
size_t offset;
@@ -117,8 +115,6 @@ See get_serialized_handle below
typedef enum triton_serialized_handle_type_t
{
MEM_PUB_HANDLE_TYPE,
- WAIT_HANDLE_TYPE /*designate the kind of handle passed to
- the wait routine waiting for acks */
}triton_serialized_handle_type_t;
@@ -126,42 +122,42 @@ typedef enum triton_serialized_handle_type_t
*Set the message size beyond which one-sided communication is used instead 2-sided for message transfers.
* The size is in bytes
*/
-triton_ret_t triton_set_one_sided_threshold(triton_msg_group_t group, size_t size);
+triton_ret_t triton_set_one_sided_threshold(triton_addr_t local_addr, size_t size);
/**
*Get the message size beyond which one-sided communication is used instead 2-sided for message transfers.
* The size is in bytes
*/
-triton_ret_t triton_get_one_sided_threshold(triton_msg_group_t group, size_t *size);
+triton_ret_t triton_get_one_sided_threshold(triton_addr_t local_addr, size_t *size);
/**
* Allocate a buffer for sending and receiving. Method can use this
* call to pin a buffer.
* TODO: Should we include information about the type of buffer to allocate?
*/
-void * triton_buffer_allocate(triton_msg_group_t group, size_t size);
+void * triton_buffer_allocate(triton_addr_t local_addr, size_t size);
/**
* Free a buffer allocated with buffer_allocate.
*/
-void triton_buffer_free(triton_msg_group_t group, void *buffer);
+void triton_buffer_free(triton_addr_t local_addr, void *buffer);
/**
* Register memory for the communication subsystem.
*/
-triton_ret_t triton_register_memory(triton_msg_group_t group, void* base,
+triton_ret_t triton_register_memory(triton_addr_t local_addr, void* base,
size_t size,
triton_mem_reg_handle_t* handle);
/**
* Unregister memory.
*/
-triton_ret_t triton_unregister_memory(triton_msg_group_t group, triton_mem_reg_handle_t handle);
+triton_ret_t triton_unregister_memory(triton_addr_t local_addr, triton_mem_reg_handle_t handle);
/**
* Expose a registered memory portion to remote peers.
*/
-__blocking triton_ret_t triton_publish_memory(triton_msg_group_t group, triton_mem_reg_handle_t in_handle,
+__blocking triton_ret_t triton_publish_memory(triton_addr_t local_addr, triton_mem_reg_handle_t in_handle,
size_t offset,
size_t size,
triton_mem_pub_handle_t* out_handle);
@@ -169,36 +165,37 @@ __blocking triton_ret_t triton_publish_memory(triton_msg_group_t group, triton_m
/**
* End access permissions to a memory portion previously exposed.
*/
-__blocking triton_ret_t triton_unpublish_memory(triton_msg_group_t group, triton_mem_pub_handle_t handle);
+__blocking triton_ret_t triton_unpublish_memory(triton_addr_t local_addr, triton_mem_pub_handle_t handle);
/*returns the size of buffer required to get the serialized representation
of the info type passed as argument
*/
-int triton_get_serialized_handle_size(triton_msg_group_t group, triton_serialized_handle_type_t handle_type);
+int triton_get_serialized_handle_size(triton_addr_t local_addr, triton_serialized_handle_type_t handle_type,
+ void* pointer_to_handle);
/*Ask the net-module to provide the serialized net-module-specific info
associated with the handle *pointer_to_in_handle.
COMMENTS:
1. pointer_to_in_handle must point to a handle compatible with handle_type
- (triton_msg_group_t group, e.g. if handle_type is MEM_PUB_HANDLE_TYPE, then pointer_to_in_handle must
+ (triton_addr_t local_addr, e.g. if handle_type is MEM_PUB_HANDLE_TYPE, then pointer_to_in_handle must
point to a triton_mem_pub_handle_t).
2. serialized_handle must be large enough. A previous call to
get_serialized_handle_size is recommended to determine the required buffer size.
- 3. This function might succeed even if set_serialized_handle (triton_msg_group_t group, below) was not called previously
+ 3. This function might succeed even if set_serialized_handle (triton_addr_t local_addr, below) was not called previously
to return pointer_to_in_handle. E.g. When publish_memory returns a handle h, any
subsequent call to get_serialized_handle with h and handle_type == MEM_PUB_HANDLE_TYPE
will succeed.
*/
-triton_ret_t triton_get_serialized_handle(triton_msg_group_t group, triton_serialized_handle_type_t handle_type,
+triton_ret_t triton_get_serialized_handle(triton_addr_t local_addr, triton_serialized_handle_type_t handle_type,
void* pointer_to_in_handle,
char* serialized_handle);
/*
- Set an serialized info and get a handle to it.
+ Create an serialized info and get a handle to it.
COMMENTS:
1. See comment 1 of get_serialized_handle above
@@ -212,30 +209,30 @@ triton_ret_t triton_get_serialized_handle(triton_msg_group_t group, triton_seria
4. Caution: 2 handles to the same serialized info might be different;
handles might not be treated by the net-module like C pointers or even reference-counted pointers.
E.g. two calls to set_serialized_handle that returned h0 and h1 must be coupled with the 2 calls
- free_serialized_handle(triton_msg_group_t group, ...,h0) and free_serialized_handle(triton_msg_group_t group, ..., h1) (triton_msg_group_t group, in any order). So, the 2 calls
- free_serialized_handle(triton_msg_group_t group, ...,h1) and free_serialized_handle(triton_msg_group_t group, ..., h1) might be wrong unless something like
- free_serialized_handle(triton_msg_group_t group, ...,h1); h1 = h0; free_serialized_handle(triton_msg_group_t group, ..., h1) happens.
+ free_serialized_handle(triton_addr_t local_addr, ...,h0) and free_serialized_handle(triton_addr_t local_addr, ..., h1) (triton_addr_t local_addr, in any order). So, the 2 calls
+ free_serialized_handle(triton_addr_t local_addr, ...,h1) and free_serialized_handle(triton_addr_t local_addr, ..., h1) might be wrong unless something like
+ free_serialized_handle(triton_addr_t local_addr, ...,h1); h1 = h0; free_serialized_handle(triton_addr_t local_addr, ..., h1) happens.
*/
-triton_ret_t triton_set_serialized_handle(triton_msg_group_t group, triton_serialized_handle_type_t handle_type,
+triton_ret_t triton_create_handle_from_serialization(triton_addr_t local_addr, triton_serialized_handle_type_t handle_type,
void* pointer_to_out_handle,
char* serialized_handle);
/*
Free an serialized info of type handle_type and whose handle is *pointer_to_in_handle.
COMMENTS:
- 1. This call will(triton_msg_group_t group, should) succeed only for handles created through set_serialized_handle.
- The function will(triton_msg_group_t group, should) fail if a handle of the same info type created through a different mechanism
- is provided. E.g. This function will(triton_msg_group_t group, should) fail for any attempt to release the info associated
+ 1. This call will(triton_addr_t local_addr, should) succeed only for handles created through set_serialized_handle.
+ The function will(triton_addr_t local_addr, should) fail if a handle of the same info type created through a different mechanism
+ is provided. E.g. This function will(triton_addr_t local_addr, should) fail for any attempt to release the info associated
with a triton_mem_pub_handle_t that was directly created through publish_memory. For such a handle
unpublish_memory is the right way to "free" things.
*/
-triton_ret_t triton_free_serialized_handle(triton_msg_group_t group, triton_serialized_handle_type_t handle_type,
+triton_ret_t triton_free_serialized_handle(triton_addr_t local_addr, triton_serialized_handle_type_t handle_type,
void* pointer_to_in_handle);
/*dest_bufs is a an array of buffer local to the initator of the operation
iov_items_count is the number of iov_items and also the number of bufs
*/
-__blocking triton_ret_t triton_msg_get(triton_msg_group_t group, triton_addr_t from,
+__blocking triton_ret_t triton_msg_get(triton_addr_t local_addr, triton_addr_t from,
triton_mem_pub_handle_t handle,
triton_iov_item_t* iovs,
char** dest_bufs,
@@ -244,19 +241,21 @@ __blocking triton_ret_t triton_msg_get(triton_msg_group_t group, triton_addr_t f
/*src_bufs is a an array of buffer local to the initator of the operation
iov_items_count is the number of iov_items and also the number of bufs
*/
-__blocking triton_ret_t triton_msg_put(triton_msg_group_t group, triton_addr_t to,
+__blocking triton_ret_t triton_msg_put(triton_addr_t local_addr, triton_addr_t to,
triton_mem_pub_handle_t handle,
triton_iov_item_t* iovs,
char** src_bufs,
int iov_items_count);
-__blocking triton_ret_t triton_start_epoch(triton_msg_group_t group, triton_mem_pub_handle_t in_handle);
+__blocking triton_ret_t triton_start_epoch(triton_addr_t local_addr, triton_mem_pub_handle_t in_handle);
-__blocking triton_ret_t triton_end_epoch(triton_msg_group_t group, triton_mem_pub_handle_t in_handle);
+__blocking triton_ret_t triton_end_epoch(triton_addr_t local_addr, triton_mem_pub_handle_t in_handle);
-__blocking triton_ret_t triton_wait(triton_msg_group_t group, triton_mem_pub_handle_t handle,
- triton_wait_handle_t wait_handle);
+/*Wait_info is netmodule-specific data that defines the wait completion condition
+*/
+__blocking triton_ret_t triton_wait(triton_addr_t local_addr, triton_mem_pub_handle_t handle,
+ void* wait_info);
#endif /* __TRITON_MESSAGE_HAE__ */
hooks/post-receive
--
1
0
24 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via f11a0cbcfe05751011104fba685afc2d44ad5f85 (commit)
from c96d8f72dcc4b8389660aa2eff5fc751fdefa5ff (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit f11a0cbcfe05751011104fba685afc2d44ad5f85
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Aug 24 12:07:12 2012 -0400
turn off block alignment in phist and block-rmw
- no need to page align any I/O operations since these tests aren't
using O_DIRECT
-----------------------------------------------------------------------
Summary of changes:
.../block-read-modify-write/block-rmw-server.ae | 2 ++
.../examples/parallel-histogram/phist-server.ae | 2 ++
2 files changed, 4 insertions(+), 0 deletions(-)
Diff of changes:
diff --git a/code/src/examples/block-read-modify-write/block-rmw-server.ae b/code/src/examples/block-read-modify-write/block-rmw-server.ae
index d162297..da36e83 100644
--- a/code/src/examples/block-read-modify-write/block-rmw-server.ae
+++ b/code/src/examples/block-read-modify-write/block-rmw-server.ae
@@ -68,6 +68,8 @@ __blocking triton_ret_t server_main (int nservers,
*/
ret = triton_zeroconf_set("triton.tosd.flags", "DATA_SYNC");
triton_error_assert(ret);
+ ret = triton_zeroconf_set("triton.tosd.alignment", "1");
+ triton_error_assert(ret);
/* set default replication factor for servers */
ret = triton_zeroconf_set("triton.rosd.default_replication", "1");
diff --git a/code/src/examples/parallel-histogram/phist-server.ae b/code/src/examples/parallel-histogram/phist-server.ae
index 56387d6..2b88e2e 100644
--- a/code/src/examples/parallel-histogram/phist-server.ae
+++ b/code/src/examples/parallel-histogram/phist-server.ae
@@ -69,6 +69,8 @@ __blocking triton_ret_t server_main (int nservers,
*/
ret = triton_zeroconf_set("triton.tosd.flags", "DATA_SYNC");
triton_error_assert(ret);
+ ret = triton_zeroconf_set("triton.tosd.alignment", "1");
+ triton_error_assert(ret);
/* set default replication factor for servers */
sprintf(rep, "%d", params->rep_factor);
hooks/post-receive
--
1
0
24 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via c96d8f72dcc4b8389660aa2eff5fc751fdefa5ff (commit)
via 0ea358186cf82382db065be36f9e4db94004923d (commit)
via 7774b2be1ec13b51538c8b3f147dbdc1783b0f38 (commit)
from ead3b1d16bd55e6b4b1bad0294c52bc32fe86f6a (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit c96d8f72dcc4b8389660aa2eff5fc751fdefa5ff
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Aug 24 11:47:02 2012 -0400
switch versioning_rmw to use generic rmw
commit 0ea358186cf82382db065be36f9e4db94004923d
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Aug 24 11:37:37 2012 -0400
update locking routine to do single block rmw
commit 7774b2be1ec13b51538c8b3f147dbdc1783b0f38
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Aug 24 11:14:36 2012 -0400
stub rmw benchmark based on phist
-----------------------------------------------------------------------
Summary of changes:
code/configure.ac | 1 +
.../README | 31 +--
.../block-rmw-client.ae} | 206 ++++----------------
.../block-rmw-server.ae} | 20 +--
.../block-rmw.ae} | 40 ++---
.../block-rmw.hae} | 9 +-
.../examples/block-read-modify-write/module.mk.in | 12 ++
7 files changed, 76 insertions(+), 243 deletions(-)
copy code/src/examples/{parallel-histogram => block-read-modify-write}/README (52%)
copy code/src/examples/{parallel-histogram/phist-client.ae => block-read-modify-write/block-rmw-client.ae} (77%)
copy code/src/examples/{parallel-histogram/phist-server.ae => block-read-modify-write/block-rmw-server.ae} (80%)
copy code/src/examples/{parallel-histogram/phist.ae => block-read-modify-write/block-rmw.ae} (86%)
copy code/src/examples/{parallel-histogram/phist.hae => block-read-modify-write/block-rmw.hae} (87%)
create mode 100644 code/src/examples/block-read-modify-write/module.mk.in
Diff of changes:
diff --git a/code/configure.ac b/code/configure.ac
index fa6e5e5..0933f1e 100644
--- a/code/configure.ac
+++ b/code/configure.ac
@@ -429,4 +429,5 @@ doc/resilience/resilience-book.txt
doc/resilience/prototype-2012-07/module.mk
src/asg/module.mk
src/examples/parallel-histogram/module.mk
+src/examples/block-read-modify-write/module.mk
])
diff --git a/code/src/examples/parallel-histogram/README b/code/src/examples/block-read-modify-write/README
similarity index 52%
copy from code/src/examples/parallel-histogram/README
copy to code/src/examples/block-read-modify-write/README
index 448d567..64ddfa4 100644
--- a/code/src/examples/parallel-histogram/README
+++ b/code/src/examples/block-read-modify-write/README
@@ -1,13 +1,11 @@
-== Parallel Histogram
+== Block read modify write
=== Introduction
-The parallel histogram example is setup as an MPI program which divides into
-clients and servers. When run, the program executes the random workload
-specified and records the number of operations over each sampling period.
-When the program exists, the report for the timing information is dumped
-to standard output.
+
+This program performs simple read/modify/write updates to portions of
+objects.
=== Build
-make src/examples/parallel-histogram/phist
+make src/examples/block-read-modify-write/block-rmw
=== Run
==== Parameters
@@ -15,25 +13,16 @@ make src/examples/parallel-histogram/phist
--nserver <n> number of MPI ranks to assign as servers, remaing become clients.
--datapoints <n> number of bin updates by each client.
--sample <n> interval in seconds to sample updates/s.
---mode <n> test mode: 0 = Locking, 1 = Conditional, 2 = Atomic
+--mode <n> test mode: 0 = Locking, 1 = Conditional
--objects <n> number objects per server, with bins divided evenly among
them. Zero means to use a separate object for each bin.
---fault-delay <n> number of seconds to wait before injecting a fault
---fault-type <n> type of fault to inject:
- 0: change status of first server to TRITON_STATUS_UNREACHABLE
- 1: inject send failures at a rate of 5% on first server
- 2: inject 100% send failures on first server, then mark it as unreachable
- 2 seconds later
---replication <n> replication factor (defaults to 1)
--timeout-ms <n> timeout in milliseconds for RPC operations (defaults to 1000)
-==== Locking / 1 bin per object
-This example runs with 1024 locks protecting the bucks and half the MPI ranks
-are clients and half are servers. There is 1 bin per object.
+==== Locking, one object per server, one fork per bin
+
+mpiexec -n 4 ./block-rmw --nservers 2 --objects 1 --datapoints 100 --mode 0 --path /tmp/rmw --sample 10 --nlocks 2 --bins 100 --zoohost localhost --concurrency 2 --dist 0 --timeout-ms 5000
-PROC=16
-TOSD=/tmp/data
-mpirun -np $PROC ./phist --nservers $(($PROC/2)) --objects 0 --datapoints 1000 --sample 10 --nlocks 1024 --mode 0 --path $TOSD
+Change the --mode argument from 0 to 1 to switch to conditional writes.
== Zookeeper
=== Installation
diff --git a/code/src/examples/parallel-histogram/phist-client.ae b/code/src/examples/block-read-modify-write/block-rmw-client.ae
similarity index 77%
copy from code/src/examples/parallel-histogram/phist-client.ae
copy to code/src/examples/block-read-modify-write/block-rmw-client.ae
index 04d5046..7afe887 100644
--- a/code/src/examples/parallel-histogram/phist-client.ae
+++ b/code/src/examples/block-read-modify-write/block-rmw-client.ae
@@ -16,7 +16,7 @@
#include "src/common/triton-error.h"
#include "src/common/triton-bootstrap.hae"
#include "src/replicated-osd/rosd-client.hae"
-#include "phist.hae"
+#include "block-rmw.hae"
#include "src/replicated-osd/rosd-utils.h"
#include "src/state/status.hae"
#include "src/state/state.hae"
@@ -30,14 +30,15 @@
#define FILE_MAX_VALUES 10000
static __blocking void locking_rmw(uint128_t oid, uint64_t fork,
- uint64_t value, int rep_factor,
+ int rep_factor,
+ char* buffer, int block_size,
int create_flag, zkr_lock_mutex_t *lock,
double *tlock, double *tupdate);
static __blocking void versioning_rmw(uint128_t oid, uint64_t fork,
- uint64_t value, int rep_factor,
+ int rep_factor,
+ char* buffer, int block_size,
int create_flag,
double *tlock, double *tupdate);
-static void atomic_update(uint128_t oid, uint64_t fork, uint64_t value);
static void zk_clean_node(zhandle_t *zh, const char *path);
static __blocking void workload (parameters_t *params,
int rank,
@@ -62,36 +63,6 @@ static void zwatcher(zhandle_t *zh,
return;
}
-/* NOTE: If we set the failure status on the actual failed server, then the
- * update won't get propagated to any other servers (because the failed
- * servers network is down). We instead send the update to a different
- * server to propagate.
- */
-static __blocking triton_ret_t status_set_other(triton_node_t failed_svr, triton_node_t status_svr)
-{
- triton_ret_t tret;
- triton_addr_t addr;
- struct state_client_update_req req;
- triton_string_t key;
- triton_string_t value;
-
- tret = triton_map_lookup_node(status_svr, &addr);
- if(tret != TRITON_SUCCESS)
- return(tret);
-
- triton_string_init(&key, "triton.status");
- triton_string_init(&value, "TRITON_STATUS_UNREACHABLE");
-
- aer_init_struct_state_client_update_req(&req, &failed_svr, &key, &value);
- tret = remote_state_client_update(AER_DEFAULT_CTX, addr, &req, NULL);
- aer_destroy_struct_state_client_update_req(&req);
-
- triton_string_destroy(&key);
- triton_string_destroy(&value);
-
- return(tret);
-}
-
__blocking triton_ret_t client_main (int nservers,
int nclients,
void *data,
@@ -119,7 +90,6 @@ __blocking triton_ret_t client_main (int nservers,
uint128_t oid;
char timeout[64];
triton_addr_t svr_addr;
- double fault_time;
zhandle_t *zk;
zkr_lock_mutex_t *zkmutex;
@@ -142,18 +112,12 @@ __blocking triton_ret_t client_main (int nservers,
ret = triton_zeroconf_set("triton.traffic_cop.static_timeout", timeout);
triton_error_assert(ret);
- /* set up fault injection network method */
- ret = triton_zeroconf_set("aesop.remote.default_net", "fault");
- triton_error_assert(ret);
- ret = triton_zeroconf_set("triton.net.fault_injector.base_method", "triton.net.mpi");
- triton_error_assert(ret);
-
ret = triton_init("triton.client");
triton_error_assert(ret);
ret = aesop_branch_threader_init();
triton_error_assert(ret);
- ret = triton_attach("fault://0");
+ ret = triton_attach("mpi://0");
triton_error_assert(ret);
if (params->mode == MODE_LOCKING)
@@ -235,7 +199,7 @@ __blocking triton_ret_t client_main (int nservers,
{
oid.l = nodes[i].l;
oid.u = nodes[i].u + j + 10000;
- ret = client_rosd_create(oid, params->rep_factor, 0);
+ ret = client_rosd_create(oid, 1, 0);
triton_error_assert(ret);
}
}
@@ -244,67 +208,9 @@ __blocking triton_ret_t client_main (int nservers,
s = MPI_Wtime();
MPI_Barrier(client_comm);
- pwait
- {
- pbranch
- {
- triton_ret_t pret;
-
- if(client_rank == 0 && params->fault_delay > 0)
- {
- pret = triton_timer(params->fault_delay * 1000);
- if(pret == TRITON_SUCCESS)
- {
- fault_time = MPI_Wtime();
- if(params->fault_type == 0)
- {
- printf("%f: Client 0 injecting failure type 0 (setting server 0 status to TRITON_STATUS_UNREACHABLE)\n", fault_time);
- pret = triton_status_set(nodes[0], TRITON_STATUS_UNREACHABLE);
- triton_error_assert(pret);
- }
- else if(params->fault_type == 1)
- {
- printf("%f: Client 0 injecting failure type 1 (injecting network send error rate of 5 percent on server 0)\n", fault_time);
- printf("WARNING: this test case is likely to fail; see #224.\n");
- pret = triton_map_lookup_node(nodes[0], &svr_addr);
- triton_error_assert(pret);
- pret = client_triton_inject_send_failures(svr_addr,
- 100, INT32_MAX, 5);
- triton_error_assert(pret);
- }
- else if(params->fault_type == 2)
- {
- printf("Client 0 injecting failure type 2 (total network failure on server 0, followed in 2 seconds by marking it with TRITON_STATUS_UNREACHABLE)\n");
- pret = triton_map_lookup_node(nodes[0], &svr_addr);
- triton_error_assert(pret);
- pret = client_triton_inject_send_failures(svr_addr,
- 100, INT32_MAX, 100);
- triton_error_assert(pret);
-
- pret = triton_timer(2000);
- triton_error_assert(pret);
- pret = status_set_other(nodes[0], nodes[1]);
- triton_error_assert(pret);
- }
- else
- {
- assert(0);
- }
- }
- else
- {
- fprintf(stderr, "WARNING: no fault injected; test did not run long enough.\n");
- }
- }
- }
- pbranch
- {
- workload(params, client_rank, nclients, nservers, &samples,
- value_count, value, zkmutex, &rates, ×, nodes,
- &timing1, &timing2);
- aesop_cancel_branches_wait();
- }
- }
+ workload(params, client_rank, nclients, nservers, &samples,
+ value_count, value, zkmutex, &rates, ×, nodes,
+ &timing1, &timing2);
MPI_Barrier(client_comm);
e = MPI_Wtime();
@@ -468,11 +374,14 @@ static __blocking void workload (parameters_t *params,
pprivate uint64_t fork;
pprivate uint64_t bin;
pprivate uint128_t oid;
+ pprivate char* buffer;
for (pb_id = 0; pb_id < params->concurrency; ++pb_id)
{
pbranch
{
+ buffer = malloc(params->block_size);
+ assert(buffer);
rc1 = pthread_mutex_lock(&mutex);
assert(rc1 == 0);
@@ -582,8 +491,8 @@ static __blocking void workload (parameters_t *params,
tlock = 0;
//printf("oid=%ld, fork=%ld, v=%ld, lock=%d lockp=%d\n",
// oid.u, fork, v, lock, lock+(pb_id*params->nlocks));
- locking_rmw(oid, fork, v,
- params->rep_factor,
+ locking_rmw(oid, fork,
+ 1, buffer, params->block_size,
params->objects_per_server,
&zkmutex[lock+(pb_id*params->nlocks)],
&tlock, &tupdate);
@@ -593,17 +502,12 @@ static __blocking void workload (parameters_t *params,
{
tupdate = 0;
tlock = 0;
- versioning_rmw(oid, fork, v,
- params->rep_factor,
+ versioning_rmw(oid, fork,
+ 1, buffer, params->block_size,
params->objects_per_server,
&tlock, &tupdate);
break;
}
- case MODE_ATOMIC:
- {
- atomic_update(oid, fork, v);
- break;
- }
}
rc1 = pthread_mutex_lock(&mutex);
@@ -667,8 +571,9 @@ static __blocking int ae_zoo_lock(zkr_lock_mutex_t *lock)
static __blocking void locking_rmw(uint128_t bin,
uint64_t subbin,
- uint64_t value,
int rep_factor,
+ char* buffer,
+ int block_size,
int create_flag,
zkr_lock_mutex_t *lock,
double *tlock,
@@ -681,19 +586,19 @@ static __blocking void locking_rmw(uint128_t bin,
uint64_t fork;
int64_t offset;
int64_t size;
- uint64_t bin_value = 0;
uint32_t flags;
int read;
int rc;
int i;
double e,s;
int rank;
+ int *value;
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
oid = bin;
fork = subbin;
offset = 0;
- size = sizeof(bin_value);
+ size = block_size;
flags = 0;
read = 0;
@@ -713,17 +618,9 @@ static __blocking void locking_rmw(uint128_t bin,
ret = client_rosd_read(oid, fork, size, offset, flags, &resp);
if (ret == TRITON_SUCCESS)
{
- if (resp.buffer.size == sizeof(bin_value))
- {
- memcpy(&bin_value, resp.buffer.buffer, sizeof(bin_value));
- }
- else if (resp.buffer.size == 0)
- {
- bin_value = 0;
- }
- else
+ if (resp.buffer.size == 0)
{
- assert(0);
+ memset(buffer, 0, block_size);
}
aer_destroy_struct_rosd_read_resp(&resp);
read = 1;
@@ -751,9 +648,10 @@ static __blocking void locking_rmw(uint128_t bin,
}
}
- bin_value += 1;
+ value = (int*)buffer;
+ *value = (*value)+1;
- triton_buffer_init(&buf, &bin_value, sizeof(bin_value));
+ triton_buffer_init(&buf, buffer, block_size);
ret = client_rosd_write(oid, fork, buf, offset, 0, flags);
if(ret != TRITON_SUCCESS)
@@ -763,16 +661,6 @@ static __blocking void locking_rmw(uint128_t bin,
assert(0);
}
- triton_buffer_init(&buf, &value, sizeof(value));
- offset = bin_value * sizeof(value);
-
- ret = client_rosd_write(oid, fork, buf, offset, 0, flags);
- if(ret != TRITON_SUCCESS)
- {
- fprintf(stderr, "client_rosd_write() failure on rank %d\n", rank);
- triton_error_print(ret, "client_rosd_write");
- assert(0);
- }
e = MPI_Wtime();
*tupdate += (e-s);
@@ -789,8 +677,9 @@ static __blocking void locking_rmw(uint128_t bin,
static __blocking void versioning_rmw(uint128_t bin,
uint64_t subbin,
- uint64_t value,
int rep_factor,
+ char* buffer,
+ int block_size,
int create_flag,
double *tlock,
double *tupdate)
@@ -810,6 +699,7 @@ static __blocking void versioning_rmw(uint128_t bin,
int rank;
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
uint64_t version;
+ int *value;
oid = bin;
fork = subbin;
@@ -820,7 +710,7 @@ static __blocking void versioning_rmw(uint128_t bin,
do
{
offset = 0;
- size = sizeof(bin_value);
+ size = block_size;
read = 0;
while(!read)
{
@@ -829,17 +719,9 @@ static __blocking void versioning_rmw(uint128_t bin,
{
version = resp.version;
version++;
- if (resp.buffer.size == sizeof(bin_value))
- {
- memcpy(&bin_value, resp.buffer.buffer, sizeof(bin_value));
- }
- else if (resp.buffer.size == 0)
+ if (resp.buffer.size == 0)
{
- bin_value = 0;
- }
- else
- {
- assert(0);
+ memset(buffer, 0, block_size);
}
aer_destroy_struct_rosd_read_resp(&resp);
read = 1;
@@ -868,9 +750,10 @@ static __blocking void versioning_rmw(uint128_t bin,
}
}
- bin_value += 1;
+ value = (int*)buffer;
+ *value = (*value)+1;
- triton_buffer_init(&buf, &bin_value, sizeof(bin_value));
+ triton_buffer_init(&buf, &buffer, block_size);
ret = client_rosd_write(oid, fork, buf, offset, version, ROSD_FLAG_COND_WRITE);
if(ret != TRITON_SUCCESS && !triton_error_equal(ret, TRITON_ERR_CONDITIONAL))
@@ -879,20 +762,6 @@ static __blocking void versioning_rmw(uint128_t bin,
triton_error_print(ret, "client_rosd_write");
assert(0);
}
-
- if(ret == TRITON_SUCCESS)
- {
- triton_buffer_init(&buf, &value, sizeof(value));
- offset = bin_value * sizeof(value);
-
- ret = client_rosd_write(oid, fork, buf, offset, 0, 0);
- if(ret != TRITON_SUCCESS)
- {
- fprintf(stderr, "client_rosd_write() failure on rank %d\n", rank);
- triton_error_print(ret, "client_rosd_write");
- assert(0);
- }
- }
}while(triton_error_equal(ret, TRITON_ERR_CONDITIONAL));
e = MPI_Wtime();
@@ -902,11 +771,6 @@ static __blocking void versioning_rmw(uint128_t bin,
return;
}
-static void atomic_update(uint128_t bin, uint64_t subbin, uint64_t value)
-{
- return;
-}
-
// delete path and child nodes
static void zk_clean_node(zhandle_t *zh, const char *path)
{
diff --git a/code/src/examples/parallel-histogram/phist-server.ae b/code/src/examples/block-read-modify-write/block-rmw-server.ae
similarity index 80%
copy from code/src/examples/parallel-histogram/phist-server.ae
copy to code/src/examples/block-read-modify-write/block-rmw-server.ae
index 56387d6..d162297 100644
--- a/code/src/examples/parallel-histogram/phist-server.ae
+++ b/code/src/examples/block-read-modify-write/block-rmw-server.ae
@@ -20,7 +20,7 @@
#include "src/state/status.hae"
#include "src/common/traffic-cop.hae"
#include "src/rebuild/rb_module.h"
-#include "phist.hae"
+#include "block-rmw.hae"
__blocking triton_ret_t server_main (int nservers,
int nclients,
@@ -32,7 +32,6 @@ __blocking triton_ret_t server_main (int nservers,
parameters_t *params;
char path[PATH_MAX];
char timeout[64];
- char rep[64];
int rc;
printf("server startup: s=%d c=%d data=%p\n", nservers, nclients, data);
@@ -71,8 +70,7 @@ __blocking triton_ret_t server_main (int nservers,
triton_error_assert(ret);
/* set default replication factor for servers */
- sprintf(rep, "%d", params->rep_factor);
- ret = triton_zeroconf_set("triton.rosd.default_replication", rep);
+ ret = triton_zeroconf_set("triton.rosd.default_replication", "1");
triton_error_assert(ret);
/* Set static timeouts */
@@ -80,20 +78,6 @@ __blocking triton_ret_t server_main (int nservers,
ret = triton_zeroconf_set("triton.traffic_cop.static_timeout", timeout);
triton_error_assert(ret);
- /* set up fault injection network method */
- ret = triton_zeroconf_set("aesop.remote.default_net", "fault");
- triton_error_assert(ret);
- ret = triton_zeroconf_set("triton.net.fault_injector.base_method", "triton.net.mpi");
- triton_error_assert(ret);
- /* TODO: Really the fakess component should use MPI directly rather than
- * the fault method (fakess itself is not intended to be fault
- * tolerant). The problem is that fakess relies on the default
- * mapping table which will be set to use fault addresses, causing
- * a mismatch if fakess tries to use MPI.
- */
- ret = triton_zeroconf_set("triton.fakess.method", "fault");
- triton_error_assert(ret);
-
ret = triton_init("triton.server");
triton_error_assert(ret);
diff --git a/code/src/examples/parallel-histogram/phist.ae b/code/src/examples/block-read-modify-write/block-rmw.ae
similarity index 86%
copy from code/src/examples/parallel-histogram/phist.ae
copy to code/src/examples/block-read-modify-write/block-rmw.ae
index 7549ebc..4ec3518 100644
--- a/code/src/examples/parallel-histogram/phist.ae
+++ b/code/src/examples/block-read-modify-write/block-rmw.ae
@@ -7,7 +7,7 @@
#include "src/aesop/aesop.h"
#include "src/common/triton-error.h"
#include "src/remote/client-server-exec.hae"
-#include "phist.hae"
+#include "block-rmw.hae"
#define POLL_TIMEOUT 1000
@@ -25,13 +25,11 @@ static void parse_args (int argc,
int *nlocks,
char *file,
int *nbins,
- int *fault_delay,
- int *fault_type,
- int *rep_factor,
char *zoohost,
int *concurrency,
int *dist,
- int *timeout_ms)
+ int *timeout_ms,
+ int *block_size)
{
static struct option long_options[] = {
{"nservers", 1, NULL, 0 },
@@ -43,13 +41,11 @@ static void parse_args (int argc,
{"nlocks", 1, NULL, 0 },
{"file", 1, NULL, 0 },
{"bins", 1, NULL, 0 },
- {"fault-delay", 1, NULL, 0},
- {"fault-type", 1, NULL, 0},
- {"replication", 1, NULL, 0},
{"zoohost", 1, NULL, 0},
{"concurrency", 1, NULL, 0},
{"dist", 1, NULL, 0},
{"timeout-ms", 1, NULL, 0},
+ {"block-size", 1, NULL, 0},
{ NULL, 0, NULL, 0 }
};
int long_index;
@@ -96,25 +92,19 @@ static void parse_args (int argc,
*nbins = atoi(optarg);
break;
case 9:
- *fault_delay = atoi(optarg);
+ strncpy(zoohost, optarg, PATH_MAX);
break;
case 10:
- *fault_type = atoi(optarg);
+ *concurrency = atoi(optarg);
break;
case 11:
- *rep_factor = atoi(optarg);
+ *dist = atoi(optarg);
break;
case 12:
- strncpy(zoohost, optarg, PATH_MAX);
+ *timeout_ms = atoi(optarg);
break;
case 13:
- *concurrency = atoi(optarg);
- break;
- case 14:
- *dist = atoi(optarg);
- break;
- case 15:
- *timeout_ms = atoi(optarg);
+ *block_size = atoi(optarg);
break;
default:
assert(0);
@@ -155,8 +145,8 @@ int main (int argc, char **argv)
main_ret = TRITON_ERR_UNKNOWN;
memset(¶ms, 0, sizeof(params));
- params.rep_factor = 1;
params.timeout_ms = 1000;
+ params.block_size = 4096;
parse_args(argc,
argv,
@@ -169,13 +159,11 @@ int main (int argc, char **argv)
¶ms.nlocks,
params.file,
¶ms.nbins,
- ¶ms.fault_delay,
- ¶ms.fault_type,
- ¶ms.rep_factor,
params.zoohost,
¶ms.concurrency,
¶ms.dist,
- ¶ms.timeout_ms);
+ ¶ms.timeout_ms,
+ ¶ms.block_size);
rc = MPI_Init_thread(&argc, &argv, MPI_THREAD_MULTIPLE, &provided);
assert(rc == MPI_SUCCESS);
@@ -192,11 +180,9 @@ int main (int argc, char **argv)
printf(" points: %d\n", params.datapoints);
printf(" mode: %d\n", params.mode);
printf(" sample: %d\n", params.sample);
- printf("fault delay: %d\n", params.fault_delay);
- printf("fault type: %d\n", params.fault_type);
- printf("replication factor: %d\n", params.rep_factor);
printf("concurrency: %d\n", params.concurrency);
printf("timeout (ms): %d\n", params.timeout_ms);
+ printf("block size (bytes): %d\n", params.block_size);
}
ae_hints_init(&hints);
diff --git a/code/src/examples/parallel-histogram/phist.hae b/code/src/examples/block-read-modify-write/block-rmw.hae
similarity index 87%
copy from code/src/examples/parallel-histogram/phist.hae
copy to code/src/examples/block-read-modify-write/block-rmw.hae
index 87d1009..9690be9 100644
--- a/code/src/examples/parallel-histogram/phist.hae
+++ b/code/src/examples/block-read-modify-write/block-rmw.hae
@@ -1,11 +1,10 @@
-#ifndef __PHIST_HAE__
-#define __PHIST_HAE__
+#ifndef __BLOCK_RMW_HAE__
+#define __BLOCK_RMW_HAE__
#include <limits.h>
#define MODE_LOCKING 0
#define MODE_VERSIONING 1
-#define MODE_ATOMIC 2
typedef struct parameters_s
{
@@ -18,12 +17,10 @@ typedef struct parameters_s
char path[PATH_MAX];
char file[PATH_MAX];
char zoohost[PATH_MAX];
- int fault_delay;
- int fault_type;
- int rep_factor;
int concurrency;
int dist;
int timeout_ms;
+ int block_size;
} parameters_t;
__blocking triton_ret_t client_main (int nservers,
diff --git a/code/src/examples/block-read-modify-write/module.mk.in b/code/src/examples/block-read-modify-write/module.mk.in
new file mode 100644
index 0000000..88ce458
--- /dev/null
+++ b/code/src/examples/block-read-modify-write/module.mk.in
@@ -0,0 +1,12 @@
+DIR := src/examples/block-read-modify-write
+
+AESOP_HDR += $(DIR)/block-rmw.hae
+
+AETESTSRC += $(DIR)/block-rmw.ae
+
+AETESTSUPPORTSRC += $(DIR)/block-rmw-server.ae
+AETESTSUPPORTSRC += $(DIR)/block-rmw-client.ae
+
+MODCFLAGS_$(DIR)/block-rmw-client := -I/usr/include/zookeeper
+MODLDFLAGS_$(DIR)/block-rmw := -lzookeeper_mt -lzoolock
+MODLINKWITH_$(DIR)/block-rmw := $(DIR)/block-rmw-server.o $(DIR)/block-rmw-client.o
hooks/post-receive
--
1
0
23 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via 7eb1189c3bac29e64bbf534ad30f6c21a428921b (commit)
from dfeffdaf4b19c48b0cdc7223b5087943cbc7ac82 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit 7eb1189c3bac29e64bbf534ad30f6c21a428921b
Author: Dries Kimpe <dkimpe(a)mcs.anl.gov>
Date: Thu Aug 23 16:32:49 2012 -0500
Update ack text
-----------------------------------------------------------------------
Summary of changes:
papers/conditional-write/paper.tex | 4 ++--
1 files changed, 2 insertions(+), 2 deletions(-)
Diff of changes:
diff --git a/papers/conditional-write/paper.tex b/papers/conditional-write/paper.tex
index cc476e3..13db756 100644
--- a/papers/conditional-write/paper.tex
+++ b/papers/conditional-write/paper.tex
@@ -409,8 +409,8 @@ Acknowledgments.
This material is based upon work supported by, or in part by
U.S. Department of Energy's Oak Ridge National Laboratory and
included the Extreme Scale Systems Center, located at ORNL and
-funded by the DoD in part by contract number 4000111689
-``Novel Software Storage Architectures''.
+funded by the DoD in part by the
+``Novel Software Storage Architectures'' contract.
\end{document}
hooks/post-receive
--
1
0
23 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via dfeffdaf4b19c48b0cdc7223b5087943cbc7ac82 (commit)
from a699a30483ed66340193cf01bc0b21196dc59806 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit dfeffdaf4b19c48b0cdc7223b5087943cbc7ac82
Author: Dries Kimpe <dkimpe(a)mcs.anl.gov>
Date: Thu Aug 23 16:06:03 2012 -0500
Add acknowledgementin case we forget
-----------------------------------------------------------------------
Summary of changes:
papers/conditional-write/paper.tex | 10 ++++++++++
1 files changed, 10 insertions(+), 0 deletions(-)
Diff of changes:
diff --git a/papers/conditional-write/paper.tex b/papers/conditional-write/paper.tex
index 1868c1c..cc476e3 100644
--- a/papers/conditional-write/paper.tex
+++ b/papers/conditional-write/paper.tex
@@ -403,6 +403,16 @@ Acknowledgments.
%derivative works, distribute copies to the public, and perform publicly
%and display publicly, by or on behalf of the Government.
+% ASG acknowledgement
+
+\newpage
+This material is based upon work supported by, or in part by
+U.S. Department of Energy's Oak Ridge National Laboratory and
+included the Extreme Scale Systems Center, located at ORNL and
+funded by the DoD in part by contract number 4000111689
+``Novel Software Storage Architectures''.
+
+
\end{document}
hooks/post-receive
--
1
0
22 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via a699a30483ed66340193cf01bc0b21196dc59806 (commit)
from 51f996387d96a84352c2e28ccd057bf854aceed8 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit a699a30483ed66340193cf01bc0b21196dc59806
Author: Cengiz Karakoyunlu <cengiz(a)lucid64-vm.mcs.anl.gov>
Date: Wed Aug 22 17:02:11 2012 -0500
Updated the Related Work section
Comments
* Updated the bibliography and paper.bib too
-----------------------------------------------------------------------
Summary of changes:
papers/conditional-write/paper.bib | 18 ++++++++++++++++++
papers/conditional-write/paper.tex | 34 +++++++++++++++++++++++++++++++++-
2 files changed, 51 insertions(+), 1 deletions(-)
Diff of changes:
diff --git a/papers/conditional-write/paper.bib b/papers/conditional-write/paper.bib
index 6975478..d57f4b6 100644
--- a/papers/conditional-write/paper.bib
+++ b/papers/conditional-write/paper.bib
@@ -280,3 +280,21 @@ doi = {http://doi.ieeecomputersociety.org/10.1109/ICPP.2008.43},
publisher = {IEEE Computer Society},
address = {Los Alamitos, CA},
}
+
+@inproceedings{OSDPDSI08,
+ title = {Revisiting the Metadata Architecture of Parallel File Systems},
+ author = {N. Ali and A. Devulapalli and D. Dalessandro and P. Wyckoff and P. Sadayappan},
+ booktitle = {Third Petascale Data Storage Workshop, Supercomputing},
+ year = {2008},
+ month = {November},
+ url = {http://www.cse.ohio-state.edu/~alin/papers/pdsw2008.pdf}
+}
+
+@inproceedings{OSDCluster08,
+ title = {An {OSD-based} Approach to Managing Directory Operations in Parallel File Systems},
+ author = {N. Ali and A. Devulapalli and D. Dalessandro and P. Wyckoff and P. Sadayappan},
+ booktitle = {IEEE International Conference on Cluster Computing},
+ year = {2008},
+ month = {September},
+ url = {http://www.cse.ohio-state.edu/~alin/papers/cluster2008.pdf}
+}
diff --git a/papers/conditional-write/paper.tex b/papers/conditional-write/paper.tex
index 94f4703..1868c1c 100644
--- a/papers/conditional-write/paper.tex
+++ b/papers/conditional-write/paper.tex
@@ -194,7 +194,39 @@ something something something else.
The related work can get out of hand. Try to limit it to some seminal
papers in other fields, and some examples of file systems that provide
-locking primitives.
+locking primitives.
+
+File systems that support locking pritimives have been an active area of research and
+Ohio Supercomputing Center's PVFS implementation ~\cite{Beowulf01} on top of object-based
+storage devices (OSDs) is an effort in this area. Initial PVFS-OSD implementation ~\cite{OSDPDSI08}
+handled the directory operations with PVFS based metadata servers rather than handling
+them directly on OSDs; because directory operations such as create, remove had to be atomic
+to ensure correctness and OSDs did not have atomic operations or locking support.
+Their follow-up study ~\cite{10.1109/SNAPI.2008.14} presents two primitives, Compare-and-Swap (CAS)
+and Fetch-and-Add (FA) to support atomic operations based on object attributes in OSDs.
+CAS atomically compares the existing value in an attribute with the given 'compare value' and if they
+are the same, it replaces the existing value with the given 'swap value'. CAS can optionally return
+the original value of the attribute. FA atomically adds the given 'add value' to the existing value
+in an attribute. Similar to CAS, FA can optionally return the original value of the attribute. Both
+CAS and FA are advisory; they let other operations on the OSD object to continue.
+Another study from OSC ~\cite{OSDCluster08} implements CAS primitive on OSDs to support
+directory operations. Directories are represented as OSD objects and directory entries are stored in
+the attributes of the OSD object by hashing directory entry name to an attribute location. They have
+two methods to perform I/O operations in a directory; locking the directory object with the CAS primitive,
+performing the operation and then releasing the lock or using the CAS primitive directly on an attribute to perform
+the operation. In the former, a well-known attribute location in the directory object is used to indicate
+the binary value of the lock status of directory object. CAS takes in a compare value and a swap value along with
+an object identifier and an attribute number. CAS compares the existing binary lock value with the compare
+value and if they are equal, it updates the attribute with the swap value and returns the original value of
+the attribute. Insert and removes first lock the directory object, perform insertion or removal and
+finally unlock the directory object again using the CAS primitive. The latter uses CAS extension to do atomic
+insert and removes in a directory. For inserts, CAS takes NULL as the compare value and object
+handle as the swap value. A successful insert operation will replace the empty attribute slot with encoded
+directory entry and return NULL as a result. Remove takes encoded directory entry as the compare
+value and NULL as the swap value. A successful remove operation will replace the given attribute slot with NULL
+and return the entry that has been removed as a result. Experimental results show that using CAS to do atomic
+attribute manipulation yields to much better performance results compared to using CAS to lock the directory
+object.
\section{Implementation strategies}
hooks/post-receive
--
1
0
22 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via 51f996387d96a84352c2e28ccd057bf854aceed8 (commit)
from a8d796408c647a9c893dd788c60b9ee6032d104b (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit 51f996387d96a84352c2e28ccd057bf854aceed8
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Wed Aug 22 14:48:39 2012 -0400
fix name
-----------------------------------------------------------------------
Summary of changes:
papers/conditional-write/paper.tex | 2 +-
1 files changed, 1 insertions(+), 1 deletions(-)
Diff of changes:
diff --git a/papers/conditional-write/paper.tex b/papers/conditional-write/paper.tex
index 2987c7c..94f4703 100644
--- a/papers/conditional-write/paper.tex
+++ b/papers/conditional-write/paper.tex
@@ -63,7 +63,7 @@ Robert Ross,\IEEEauthorrefmark{1}
Lee Ward,\IEEEauthorrefmark{2}
Matthew Curry,\IEEEauthorrefmark{2} \\
Ruth Klundt,\IEEEauthorrefmark{2}
-Geoff Danielson,\IEEEauthorrefmark{2}
+Geoffrey Danielson,\IEEEauthorrefmark{2}
Cengiz Karakoyunlu,\IEEEauthorrefmark{3}
John Chandy,\IEEEauthorrefmark{3}
Bradley Settlemyer,\IEEEauthorrefmark{4}} \\
hooks/post-receive
--
1
0
22 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via a8d796408c647a9c893dd788c60b9ee6032d104b (commit)
from 13a8bf2057e781b63c4fcdd7a3ad6f3557285b89 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit a8d796408c647a9c893dd788c60b9ee6032d104b
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Wed Aug 22 14:45:26 2012 -0400
misc edits
-----------------------------------------------------------------------
Summary of changes:
papers/conditional-write/paper.tex | 34 +++++++++++++++++-----------------
1 files changed, 17 insertions(+), 17 deletions(-)
Diff of changes:
diff --git a/papers/conditional-write/paper.tex b/papers/conditional-write/paper.tex
index eb63a03..2987c7c 100644
--- a/papers/conditional-write/paper.tex
+++ b/papers/conditional-write/paper.tex
@@ -128,7 +128,7 @@ of access coordination is required to insure that concurrent updates do not
produce incoherent results. In some
cases this coordination is handled by the application itself (e.g. via
MPI synchronization or explicit data partitioning), while in other cases
-coordination is provided explicitly or implicitly by the storage system.
+coordination is provided by the storage system.
Storage system synchronization primitives are especially critical
in cases where it is difficult for applications to protect
consistency on their own. Examples include:
@@ -136,7 +136,7 @@ consistency on their own. Examples include:
\begin{itemize}
\item concurrent name space (file or directory) updates
\item unaligned access in block-based storage systems
-\item inherently uncoordinated updates applications themselves
+\item inherently uncoordinated writes from applications themselves
(examples here or elsewhere in text,
like reverse index computation, gups, parallel histogram, the reduction
stage of mapreduce, data sieving, etc.)
@@ -144,7 +144,7 @@ consistency on their own. Examples include:
The traditional approach to providing access coordination in these cases is via
distributed file locking in a parallel file system. There are drawbacks
-to this approach however. Distributed locking is pessimistic; it
+to this approach, however. Distributed locking is pessimistic; it
introduces significant overhead even in cases where conflicts are rare.
In addition, it introduces shared state on client processes. This shared
state complicates fault tolerance and compromizes scalability by forcing the
@@ -154,14 +154,15 @@ This problem has been studied extensively in database and shared memory
literature, and a well-known alternative to locking in those arenas is
to leverage conditional storage operations, such as compare-and-swap
or load-link/store-conditional. These primitives can be used as building
-blocks to construct a number of higher-level atomic operations.
-These conditional storage primitives have not been adopted
+blocks to construct a number of higher-level atomic operations. Conditional
+storage operations have not been adopted
in HPC storage systems to date, however, for a number of reasons. Notably,
-the traditional POSIX API does not support store conditional primitives, and
-most present-day HPC storage systems do not expose any other API to client
+the traditional POSIX API does not support conditional storage primitives, and
+present-day HPC storage systems do not expose any other API to client
nodes.
-In addition, there is a lack of support for efficient
-comparison or conditional checks in most local storage abstractions as well
+In addition, there is also a lack of support for efficient
+comparison or conditional checks in the local storage abstractions that form
+the foundation of most HPC storage systems
(meaning things like Trove (PVFS), Objectstore (Ceph), and ldiskfs or ZFS
(Lustre).
@@ -173,21 +174,20 @@ low level model for storage access is the distributed object storage
model (cite osc/pvfs, rados, s3, etc.). Such object storage models are
not tied to legacy POSIX semantics and can more easily be modified
to expose additional synchronization primitives to applications
-while still providing a solid foundation for a variety of high level
+while still providing a solid foundation for a variety of higher level
interfaces (including POSIX). The second trend is towards increased
functionality in local storage abstractions in order to support features
such as checksumming, log-structured storage, provenance, deduplication,
and versioning. These features inherently require additional metadata
support in the local storage component of a distributed storage system,
-and this same infrastructure can often be reused to support conditional
-storage primitives with modest infrastructure modification.
+and this same metadata infrastructure can often be reused to support conditional
+storage primitives with modest incremental complexity.
In this work we propose a model for providing efficient store-conditional
operators in a distributed storage system. We implement this model in a
-prototype distributed object storage system. We then evaluate the
+prototype system and evaluate the
performance of an uncoordinated read/modify/write workload using both
-store-conditional operations and traditional locking techniques in order
-study the performance tradeoffs between the two systems. We show that
+store-conditional operations and traditional locking techniques. We show that
something something something else.
\section{Related Work}
@@ -200,7 +200,7 @@ locking primitives.
The first consideration in implementing conditional operators is to choose a
comparison mechanism.
-The simplest approach is to provide a compare-and-swap
+The simplest approach is to use a compare-and-swap
primitive. This requires no modification to the read() operation,
and requires no additional metadata in the local storage abstraction. One
drawback is that it requires reading the entire affected region before
@@ -248,7 +248,7 @@ case the server can simply perform the write and rely on the local storage
abstraction to use local transactions on its storage index to perform the
comparison. The latter is likely to be more efficient, but it has a
significant drawback in that it may produce non-determistic results if the
-conditional is applied to a replicated object if all replicas do note
+conditional is applied to a replicated object and replicas do not
recieve I/O requests in the same order.
We therefore elected to use the server-based
comparison approach, because one server can be chosen to perform the
hooks/post-receive
--
1
0
22 Aug '12
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "".
The branch, master has been updated
via 13a8bf2057e781b63c4fcdd7a3ad6f3557285b89 (commit)
from 913fd187177779e6b5b532b036fc1fa8170454a9 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
commit 13a8bf2057e781b63c4fcdd7a3ad6f3557285b89
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Wed Aug 22 13:44:32 2012 -0400
add authors
-----------------------------------------------------------------------
Summary of changes:
papers/conditional-write/paper.tex | 17 +++++++++++++++++
1 files changed, 17 insertions(+), 0 deletions(-)
Diff of changes:
diff --git a/papers/conditional-write/paper.tex b/papers/conditional-write/paper.tex
index 42ced9d..eb63a03 100644
--- a/papers/conditional-write/paper.tex
+++ b/papers/conditional-write/paper.tex
@@ -56,6 +56,23 @@
%\title{VOSD Paper Draft}
\title{A Title}
+\author{\IEEEauthorblockN{Philip Carns,\IEEEauthorrefmark{1}
+Kevin Harms,\IEEEauthorrefmark{1}
+Dries Kimpe,\IEEEauthorrefmark{1}
+Robert Ross,\IEEEauthorrefmark{1}
+Lee Ward,\IEEEauthorrefmark{2}
+Matthew Curry,\IEEEauthorrefmark{2} \\
+Ruth Klundt,\IEEEauthorrefmark{2}
+Geoff Danielson,\IEEEauthorrefmark{2}
+Cengiz Karakoyunlu,\IEEEauthorrefmark{3}
+John Chandy,\IEEEauthorrefmark{3}
+Bradley Settlemyer,\IEEEauthorrefmark{4}} \\
+\IEEEauthorblockA{\IEEEauthorrefmark{1}Argonne National Laboratory, Argonne, IL, USA}
+\IEEEauthorblockA{\IEEEauthorrefmark{2}Sandia National Laboratories, Livermore, CA, USA}
+\IEEEauthorblockA{\IEEEauthorrefmark{3}University of Connecticut, Storrs, CT, USA}
+\IEEEauthorblockA{\IEEEauthorrefmark{3}Oak Ridge National Laboratory, Oak Ridge, TN, USA}
+}
+
%\author{
%{Philip Carns, Robert Ross, and Samuel Lang}
%\vspace{1.6mm}\\
hooks/post-receive
--
1
0