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
April 2014
- 1 participants
- 46 discussions
27 Apr '14
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 056effdf2631aeb0bbc85e11e7c5a97920c14643 (commit)
via 5d507aa529c7c57961b10283e3e9003e7a8eef92 (commit)
via 1f4b68d19d951a4732a32eea456b22f28f31c4f1 (commit)
via 948907e0c766b0f5e1f4bad1ea9de33bf2de19b7 (commit)
via f59a367b2b8334fb28318b42ffe3ee8e483ce8b4 (commit)
via 7e355cb8968e2e785867e702ee9073d23ea14a32 (commit)
via bbdbb1229d4281409040666ba0489486d809d6d0 (commit)
via d0ecd3fbaf8c5fe3df48672b7482016870f1df0d (commit)
via 8df78189eec3cce18fea516609f409e4345812dd (commit)
via 8860feb91cafe6f6be8ee166c33d4eec7da62503 (commit)
via 7862829659b1f8844f2e058aad99d85f38c2bfcf (commit)
via f5385cd8fd088bf4a146c823f31e38cb559dd05e (commit)
via 5d8e5aaa8910c3ad52daf9be5c61a27cbc955003 (commit)
via 7842dc0aa74f5d6dc9a8ec5fc25a3f32c8067ba6 (commit)
via 3cd538544e5c8e223ea35991ae4fd39e298a1e6a (commit)
via c8118f1185e5bc26a8252131839bb1559c7ef8d8 (commit)
via d6d42eb12f4c3eee62758701a9a948afe1b82688 (commit)
via ba7178daa0efdea03245ebca1707fc68e59e1e19 (commit)
via 5b111a97c8df0ccab0287c94337f89004602c9e7 (commit)
via ecfda6380e6119eb1c7739e7f54fd5795d24d175 (commit)
via 035ae077293f5c81c31eecd028d22ce5d7561fab (commit)
via 59d9c8b3159722ba5d9e1b3a4a04ebc9733c55c1 (commit)
via 2951c7427e96f3caedfb51e2ac89b118f7770d62 (commit)
via d1fd10f7a6caea9cd6ff47f98bdd286b2dd72ee2 (commit)
via 58dce24dbe05a8fe29c05db3be83386e96bf880f (commit)
from c347c633e3ae197603ebd7ae3428935efb2d8491 (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 056effdf2631aeb0bbc85e11e7c5a97920c14643
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 21:40:05 2014 -0400
explicit rule to gen asg-internal.h before asg.o
commit 5d507aa529c7c57961b10283e3e9003e7a8eef92
Merge: c347c633e3ae197603ebd7ae3428935efb2d8491 1f4b68d19d951a4732a32eea456b22f28f31c4f1
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 21:36:01 2014 -0400
Merge remote-tracking branch 'origin/trac-277-bulkio'
Merging complete write and read pipelining to master
-----------------------------------------------------------------------
Summary of changes:
code/src/admin-tools/triton-cp.ae | 2 +-
code/src/asg/Makefile.subdir | 2 +
code/src/replicated-osd/buffer-mgmt.ae | 6 +-
code/src/replicated-osd/buffer-mgmt.hae | 2 +-
code/src/replicated-osd/rosd-read.ae | 114 +++-------
code/src/replicated-osd/rosd-write.ae | 261 ++++++++++++++++++-----
code/src/replicated-osd/rosd.ae | 77 +++++++
code/src/replicated-osd/rosd.hae | 14 ++
code/src/transactional-osd/transactional-osd.ae | 5 +-
code/tests/Makefile.subdir | 1 +
code/tests/triton-cp-big.sh | 70 ++++++
code/tests/triton-cp.sh | 11 +-
12 files changed, 421 insertions(+), 144 deletions(-)
create mode 100755 code/tests/triton-cp-big.sh
Diff of changes:
diff --git a/code/src/admin-tools/triton-cp.ae b/code/src/admin-tools/triton-cp.ae
index a27498b..97c812a 100644
--- a/code/src/admin-tools/triton-cp.ae
+++ b/code/src/admin-tools/triton-cp.ae
@@ -17,7 +17,7 @@
#include "src/replicated-osd/rosd.hae"
/* TODO: make this configurable */
-#define BUFFER_SZ (4*1024*1024)
+#define BUFFER_SZ (32*1024*1024)
enum obj_ref_type
{
diff --git a/code/src/asg/Makefile.subdir b/code/src/asg/Makefile.subdir
index 4e5e677..b31ffd0 100644
--- a/code/src/asg/Makefile.subdir
+++ b/code/src/asg/Makefile.subdir
@@ -8,3 +8,5 @@ AE_SRC += \
AE_HDR += \
src/asg/asg-internal.hae
+
+src/asg/asg.o: src/asg/asg-internal.h
diff --git a/code/src/replicated-osd/buffer-mgmt.ae b/code/src/replicated-osd/buffer-mgmt.ae
index c5c8e91..c75ca09 100644
--- a/code/src/replicated-osd/buffer-mgmt.ae
+++ b/code/src/replicated-osd/buffer-mgmt.ae
@@ -15,7 +15,7 @@ struct buffer_mgmt_instance
struct buffer_mgmt_token
{
char* buffer;
- int size;
+ int64_t size;
struct buffer_mgmt_instance* instance;
};
@@ -72,10 +72,10 @@ struct buffer_mgmt_instance* buffer_mgmt_create_instance(
__blocking triton_ret_t buffer_mgmt_alloc(
struct buffer_mgmt_instance *instance,
- int requested_size, int* allocated_size,
+ int64_t requested_size, int64_t* allocated_size,
char** buffer, struct buffer_mgmt_token** token)
{
- int size_to_alloc = 0;
+ int64_t size_to_alloc = 0;
int ret;
assert(requested_size > 0);
diff --git a/code/src/replicated-osd/buffer-mgmt.hae b/code/src/replicated-osd/buffer-mgmt.hae
index 88fbc2f..d887ab1 100644
--- a/code/src/replicated-osd/buffer-mgmt.hae
+++ b/code/src/replicated-osd/buffer-mgmt.hae
@@ -14,7 +14,7 @@ struct buffer_mgmt_instance* buffer_mgmt_create_instance(
__blocking triton_ret_t buffer_mgmt_alloc(
struct buffer_mgmt_instance* instance,
- int requested_size, int* allocated_size,
+ int64_t requested_size, int64_t* allocated_size,
char** buffer, struct buffer_mgmt_token** token);
void buffer_mgmt_free(struct buffer_mgmt_token* token);
diff --git a/code/src/replicated-osd/rosd-read.ae b/code/src/replicated-osd/rosd-read.ae
index d4c4b1e..e0f05c7 100644
--- a/code/src/replicated-osd/rosd-read.ae
+++ b/code/src/replicated-osd/rosd-read.ae
@@ -124,68 +124,6 @@ __blocking triton_ret_t remote_triton_rpc_rosd_read(
}
-static __blocking triton_ret_t rosd_read_setup_one_buffer(
- aesop_sem_t *sem,
- int* size_remaining,
- int64_t *remote_offset,
- int64_t *local_offset,
- int* this_size,
- int64_t* this_local_offset,
- int64_t* this_remote_offset,
- struct buffer_mgmt_token **token,
- char** buffer
- )
-{
- int ret;
- triton_ret_t tret;
-
- ret = aesop_sem_down(sem);
- if(ret != AE_SUCCESS)
- {
- if(ret == AE_ERR_CANCELLED)
- return(TRITON_ERR_CANCELED);
- else
- return(TRITON_ERR_UNKNOWN);
- }
-
- if((*size_remaining) == 0)
- {
- *this_size = 0; /* done */
- ret = aesop_sem_up(sem);
- assert(ret == AE_SUCCESS);
- return(TRITON_SUCCESS);
- }
-
- /* do we have a buffer for this pbranch yet? */
- if((*buffer) == NULL)
- {
- tret = buffer_mgmt_alloc(rosd_read_buffers,
- *size_remaining, this_size, buffer, token);
- if(triton_is_error(tret))
- {
- ret = aesop_sem_up(sem);
- assert(ret == AE_SUCCESS);
- return(tret);
- }
- }
-
- /* figure out what region we are accessing now */
- if(*this_size > *size_remaining)
- *this_size = *size_remaining;
- *this_local_offset = *local_offset;
- *this_remote_offset = *remote_offset;
-
- /* update offsets/sizes for overall xfer */
- *size_remaining -= *this_size;
- *local_offset += *this_size;
- *remote_offset += *this_size;
-
- ret = aesop_sem_up(sem);
- assert(ret == AE_SUCCESS);
-
- return(TRITON_SUCCESS);
-}
-
static __blocking triton_ret_t rosd_read_xfer_one_buffer(
na_addr_t src_addr,
char* tmp_buffer,
@@ -262,8 +200,7 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
int my_position;
int ret;
na_addr_t src_addr;
- /* TODO: 64 bit? */
- int size_remaining = 0;
+ int64_t size_remaining = 0;
int64_t remote_offset = 0;
int64_t local_offset = 0;
aesop_sem_t sem;
@@ -306,8 +243,7 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
pwait
{
pprivate int i = 0;
- /* TODO: 64 bit? */
- pprivate int this_size = 0;
+ pprivate int64_t this_size = 0;
pprivate int64_t this_out_size = 0;
pprivate int64_t this_local_offset = 0;
pprivate int64_t this_remote_offset = 0;
@@ -328,7 +264,8 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
* using a semaphore
*/
//triton_debug(triton_dbg_rosd, "triton_rpc_rosd_read() pbranch %d calling setup_one_buffer().\n", i);
- tret = rosd_read_setup_one_buffer(
+ tret = rosd_pipeline_setup_one_buffer(
+ rosd_read_buffers,
&sem,
&size_remaining,
&remote_offset,
@@ -338,8 +275,18 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
&this_remote_offset,
&token,
&buffer);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
+ triton_mutex_lock(&output_mutex);
+ if(triton_is_error(tret) && !triton_is_error(out.tret))
+ {
+ out.tret = tret;
+ this_size = 0;
+ }
+ else if(triton_is_error(tret))
+ {
+ triton_error_destroy(tret);
+ }
+ triton_mutex_unlock(&output_mutex);
+
//triton_debug(triton_dbg_rosd, "triton_rpc_rosd_read() pbranch %d finished setup_one_buffer(), size_remaining: %d.\n", i, size_remaining);
if(this_size > 0)
@@ -357,21 +304,30 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
in.bulk_handle,
&this_out_size,
&this_txn_number);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
//triton_debug(triton_dbg_rosd, "triton_rpc_rosd_read() pbranch %d finished xfer_one_buffer().\n", i);
/* accumulate results in rpc response */
triton_mutex_lock(&output_mutex);
- out.out_size += this_out_size;
- if(out.txn_number == 0 && out.txn_number == this_txn_number)
- out.txn_number = this_txn_number;
+ if(triton_is_error(tret) && !triton_is_error(out.tret))
+ {
+ out.tret = tret;
+ }
+ else if(triton_is_error(tret))
+ {
+ triton_error_destroy(tret);
+ }
else
{
- /* TODO: what are we supposed to do on mixed txn
- * numbers?
- */
- out.txn_number = 0;
+ out.out_size += this_out_size;
+ if(out.txn_number == 0 && out.txn_number == this_txn_number)
+ out.txn_number = this_txn_number;
+ else
+ {
+ /* TODO: what are we supposed to do on mixed txn
+ * numbers?
+ */
+ out.txn_number = 0;
+ }
}
triton_mutex_unlock(&output_mutex);
}
@@ -379,7 +335,7 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
{
this_out_size = this_size;
}
- }while(this_out_size > 0);
+ }while(this_out_size > 0 && !triton_is_error(out.tret));
if(token != NULL)
{
diff --git a/code/src/replicated-osd/rosd-write.ae b/code/src/replicated-osd/rosd-write.ae
index f0efe74..776788d 100644
--- a/code/src/replicated-osd/rosd-write.ae
+++ b/code/src/replicated-osd/rosd-write.ae
@@ -17,6 +17,12 @@
#include "src/system-state/system-state.hae"
#include "src/transactional-osd/transactional-osd.hae"
+/* TODO: make this configurable */
+/* this is the number of concurrent buffers/transfers that the ROSD will
+ * keep in flight for a single write operation
+ */
+#define ROSD_WRITE_XFER_PIPELINE_DEPTH 4
+
/* Mercury RPC structures for rosd_write */
MERCURY_GEN_PROC(triton_rpc_rosd_write_out_t, ((triton_ret_t)(tret)))
MERCURY_GEN_PROC(triton_rpc_rosd_write_in_t,
@@ -31,6 +37,7 @@ MERCURY_GEN_PROC(triton_rpc_rosd_write_in_t,
((hg_bulk_t)(bulk_handle)))
+extern struct buffer_mgmt_instance* rosd_write_buffers;
static hg_id_t rpc_rosd_write_id;
static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle);
static int triton_rpc_rosd_write_handler(hg_handle_t handle);
@@ -96,7 +103,202 @@ __blocking triton_ret_t remote_triton_rpc_rosd_write(
}
-/* TODO: refactor this function; there is a lot going on in here now */
+static __blocking triton_ret_t rosd_write_pull_one_buffer(
+ na_addr_t src_addr,
+ char* tmp_buffer,
+ uint128_t oid,
+ uint64_t fork,
+ int64_t size,
+ int64_t remote_offset,
+ hg_bulk_t remote_bulk_handle)
+{
+ triton_ret_t tret;
+ hg_bulk_t bulk_handle = HG_BULK_NULL;
+ hg_bulk_request_t bulk_request;
+ int ret;
+
+ ret = HG_Bulk_handle_create(tmp_buffer, size, HG_BULK_READWRITE, &bulk_handle);
+ if(ret != HG_SUCCESS)
+ {
+ triton_error_msg("HG_Bulk_handle_create() failure.\n");
+ return(TRITON_ERR_NOMEM);
+ }
+
+ ret = HG_Bulk_read(src_addr, remote_bulk_handle, remote_offset,
+ bulk_handle, 0, size, &bulk_request);
+ if(ret != HG_SUCCESS)
+ {
+ triton_error_msg("HG_Bulk_read() failure.\n");
+ HG_Bulk_handle_free(bulk_handle);
+ return(TRITON_ERR_UNKNOWN);
+ }
+
+ tret = triton_mercury_bulk_wait(bulk_request);
+ if(triton_is_error(tret))
+ {
+ triton_error_msg("triton_mercury_bulk_wait() failure.\n");
+ }
+
+ HG_Bulk_handle_free(bulk_handle);
+
+ return(tret);
+}
+
+static __blocking triton_ret_t rosd_write_run_iteration(
+ na_addr_t src_addr,
+ na_addr_t next_addr,
+ int my_position,
+ int from_client_flag,
+ triton_rpc_rosd_write_in_t* in,
+ aesop_sem_t *sem,
+ int64_t* size_remaining,
+ int64_t* this_size,
+ int64_t *remote_offset,
+ int64_t *local_offset,
+ struct buffer_mgmt_token **token,
+ char** buffer,
+ int* more_flag)
+{
+ int64_t this_local_offset = 0;
+ int64_t this_remote_offset = 0;
+ triton_ret_t tret;
+
+ *more_flag = 0;
+
+ /* calculate how much to transfer in this step */
+
+ /* note that this function protects shared variables
+ * using a semaphore
+ */
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling setup_one_buffer().\n");
+ tret = rosd_pipeline_setup_one_buffer(
+ rosd_write_buffers,
+ sem,
+ size_remaining,
+ remote_offset,
+ local_offset,
+ this_size,
+ &this_local_offset,
+ &this_remote_offset,
+ token,
+ buffer);
+ if(triton_is_error(tret))
+ {
+ return(tret);
+ }
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished setup_one_buffer(), this_size: %d, this_local_offset: %ld, this_remote_offset %ld.\n", *this_size, this_local_offset, this_remote_offset);
+
+ if(*this_size > 0)
+ {
+ /* perform buffer transfer */
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling pull_one_buffer().\n");
+ tret = rosd_write_pull_one_buffer(
+ src_addr,
+ *buffer,
+ in->oid,
+ in->oid_fork,
+ *this_size,
+ this_remote_offset,
+ in->bulk_handle);
+ if(triton_is_error(tret))
+ {
+ return(tret);
+ }
+
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished pull_one_buffer().\n");
+
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch starting do_work().\n");
+ tret = rosd_write_do_work(next_addr, in->oid,
+ in->oid_fork, *buffer, *this_size, this_local_offset,
+ in->flags, in->replication_factor, in->txn_number,
+ my_position, from_client_flag);
+ if(triton_is_error(tret))
+ {
+ return(tret);
+ }
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished do_work().\n");
+
+ *more_flag = 1;
+ }
+
+ return(TRITON_SUCCESS);
+}
+
+static __blocking triton_ret_t rosd_write_run_pipeline(
+ na_addr_t src_addr,
+ na_addr_t next_addr,
+ int my_position,
+ int from_client_flag,
+ triton_rpc_rosd_write_in_t* in
+)
+{
+ int64_t size_remaining = 0;
+ int64_t remote_offset = 0;
+ int64_t local_offset = 0;
+ aesop_sem_t sem;
+ triton_ret_t out_tret = TRITON_SUCCESS;
+ triton_mutex_t output_mutex;
+
+ size_remaining = in->size;
+ local_offset = in->offset;
+ remote_offset = 0;
+ aesop_sem_init(&sem, 1);
+
+ triton_mutex_init(&output_mutex, NULL);
+
+ pwait
+ {
+ pprivate int i = 0;
+ pprivate triton_ret_t tret;
+ pprivate struct buffer_mgmt_token *token = NULL;
+ pprivate char* buffer = NULL;
+ pprivate int more_flag = 1;
+ pprivate int64_t this_size = 0;
+
+ //triton_debug(triton_dbg_rosd,
+ // "rosd_write_run_pipeline() starting pwait, size_remaining: %d.\n",
+ // size_remaining);
+ for(i=0; i<ROSD_WRITE_XFER_PIPELINE_DEPTH; i++)
+ {
+ pbranch
+ {
+ //triton_debug(triton_dbg_rosd,
+ // "rosd_write_run_pipeline() starting pbranch %d.\n", i);
+
+ do
+ {
+ tret = rosd_write_run_iteration(src_addr, next_addr,
+ my_position, from_client_flag, in, &sem,
+ &size_remaining, &this_size, &remote_offset,
+ &local_offset, &token, &buffer, &more_flag);
+ }while(more_flag && !triton_is_error(tret));
+
+ if(token != NULL)
+ {
+ buffer_mgmt_free(token);
+ }
+ //triton_debug(triton_dbg_rosd,
+ // "rosd_write_run_pipeline() finishing pbranch %d.\n", i);
+ triton_mutex_lock(&output_mutex);
+ if(triton_is_error(tret) && !triton_is_error(out_tret))
+ {
+ out_tret = tret;
+ }
+ else if(triton_is_error(tret))
+ {
+ triton_error_destroy(tret);
+ }
+ triton_mutex_unlock(&output_mutex);
+ }
+ }
+ }
+ aesop_sem_destroy(&sem);
+ triton_mutex_destroy(&output_mutex);
+
+ return(out_tret);
+}
+
+
static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
{
triton_rpc_rosd_write_out_t out;
@@ -105,16 +307,13 @@ static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
na_addr_t next_addr;
na_addr_t* addr_array;
int my_position;
+ int from_client_flag = 0;
int64_t obj_offset;
int64_t size;
int64_t out_size;
char* buffer_offsets[1];
- int from_client_flag = 0;
struct txn_nr_cache_entry* entry_p = NULL;
- char* tmp_buffer;
- hg_bulk_t bulk_handle = HG_BULK_NULL;
na_addr_t src_addr;
- hg_bulk_request_t bulk_request;
triton_debug(triton_dbg_rosd, "Called triton_rpc_rosd_write()\n");
@@ -218,61 +417,13 @@ static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
next_addr = addr_array[my_position+1];
}
- /* TODO: buffer management */
- /* TODO: pipelining */
- if(in.size > 0)
- {
- tmp_buffer = malloc(in.size);
- if(!tmp_buffer)
- {
- out.tret = TRITON_ERR_NOMEM;
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
- }
-
- ret = HG_Bulk_handle_create(tmp_buffer, in.size, HG_BULK_READWRITE, &bulk_handle);
- if(ret != HG_SUCCESS)
- {
- triton_error_msg("HG_Bulk_handle_create() failure.\n");
- out.tret = TRITON_ERR_NOMEM;
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
- }
-
- ret = HG_Bulk_read(src_addr, in.bulk_handle, 0, bulk_handle,
- 0, in.size, &bulk_request);
- if(ret != HG_SUCCESS)
- {
- triton_error_msg("HG_Bulk_write() failure.\n");
- out.tret = TRITON_ERR_UNKNOWN;
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
- }
-
- out.tret = triton_mercury_bulk_wait(bulk_request);
- if(triton_is_error(out.tret))
- {
- triton_error_msg("triton_mercury_bulk_wait() failure.\n");
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
- }
- }
-
- out.tret = rosd_write_do_work(next_addr, in.oid, in.oid_fork,
- tmp_buffer, in.size, in.offset, in.flags,
- in.replication_factor, in.txn_number, my_position, from_client_flag);
+ out.tret = rosd_write_run_pipeline(src_addr, next_addr, my_position,
+ from_client_flag, &in);
if(entry_p)
txn_nr_cache_put(entry_p);
triton_mercury_start_output(handle, &out);
- if(in.size > 0)
- HG_Bulk_handle_free(bulk_handle);
- free(tmp_buffer);
free(addr_array);
return(TRITON_SUCCESS);
diff --git a/code/src/replicated-osd/rosd.ae b/code/src/replicated-osd/rosd.ae
index 4017aba..30424b6 100644
--- a/code/src/replicated-osd/rosd.ae
+++ b/code/src/replicated-osd/rosd.ae
@@ -385,6 +385,83 @@ __blocking void trigger_server_fault(triton_ret_t tret)
return;
}
+__blocking triton_ret_t rosd_pipeline_setup_one_buffer(
+ struct buffer_mgmt_instance* buffer_instance,
+ aesop_sem_t *sem,
+ int64_t* size_remaining,
+ int64_t *remote_offset,
+ int64_t *local_offset,
+ int64_t* this_size,
+ int64_t* this_local_offset,
+ int64_t* this_remote_offset,
+ struct buffer_mgmt_token **token,
+ char** buffer
+ )
+{
+ int ret;
+ triton_ret_t tret;
+
+ ret = aesop_sem_down(sem);
+ if(ret != AE_SUCCESS)
+ {
+ if(ret == AE_ERR_CANCELLED)
+ return(TRITON_ERR_CANCELED);
+ else
+ return(TRITON_ERR_UNKNOWN);
+ }
+
+ triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() acquired semaphore with size_remaining %lld\n", lld(*size_remaining));
+
+ if((*size_remaining) == 0)
+ {
+ *this_size = 0; /* done */
+ *this_local_offset = 0;
+ *this_remote_offset = 0;
+ ret = aesop_sem_up(sem);
+ assert(ret == AE_SUCCESS);
+ return(TRITON_SUCCESS);
+ }
+
+ /* do we have a buffer for this pbranch yet? */
+ if((*buffer) == NULL)
+ {
+ tret = buffer_mgmt_alloc(buffer_instance,
+ *size_remaining, this_size, buffer, token);
+ if(triton_is_error(tret))
+ {
+ ret = aesop_sem_up(sem);
+ assert(ret == AE_SUCCESS);
+ return(tret);
+ }
+ }
+ else
+ {
+ /* If we have a buffer, then it better be non-zero size and have a
+ * token associated with it
+ */
+ assert(*this_size > 0);
+ assert(*token);
+ }
+
+ /* figure out what region we are accessing now */
+ if(*this_size > *size_remaining)
+ *this_size = *size_remaining;
+ *this_local_offset = *local_offset;
+ *this_remote_offset = *remote_offset;
+
+ /* update offsets/sizes for overall xfer */
+ *size_remaining -= *this_size;
+ *local_offset += *this_size;
+ *remote_offset += *this_size;
+
+ triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() releasing semaphore with size_remaining %lld\n", lld(*size_remaining));
+
+ ret = aesop_sem_up(sem);
+ assert(ret == AE_SUCCESS);
+
+ return(TRITON_SUCCESS);
+}
+
/*
* Local Variables:
* c-basic-offset: 4
diff --git a/code/src/replicated-osd/rosd.hae b/code/src/replicated-osd/rosd.hae
index 44d3c17..633fea8 100644
--- a/code/src/replicated-osd/rosd.hae
+++ b/code/src/replicated-osd/rosd.hae
@@ -2,10 +2,12 @@
#define __ROSD_HAE__
#include <aesop/aesop.h>
+#include <aesop/sem.hae>
#include <mercury.h>
#include <triton-uint128.h>
#include "src/common/triton-error.h"
+#include "src/replicated-osd/buffer-mgmt.hae"
/* for requests that should be fanned out from the master for replication */
#define ROSD_FLAG_FANOUT 1
@@ -62,6 +64,18 @@ __blocking triton_ret_t remote_triton_rpc_rosd_read(
int64_t* out_size,
uint32_t flags);
+__blocking triton_ret_t rosd_pipeline_setup_one_buffer(
+ struct buffer_mgmt_instance* buffer_instance,
+ aesop_sem_t *sem,
+ int64_t* size_remaining,
+ int64_t *remote_offset,
+ int64_t *local_offset,
+ int64_t* this_size,
+ int64_t* this_local_offset,
+ int64_t* this_remote_offset,
+ struct buffer_mgmt_token **token,
+ char** buffer
+ );
#endif /* __ROSD_HAE */
diff --git a/code/src/transactional-osd/transactional-osd.ae b/code/src/transactional-osd/transactional-osd.ae
index 0d7b900..65aa8c8 100644
--- a/code/src/transactional-osd/transactional-osd.ae
+++ b/code/src/transactional-osd/transactional-osd.ae
@@ -424,9 +424,10 @@ static __blocking triton_ret_t end_write(struct coalesce_obj *obj)
}
/* notify everyone else in group */
- for(i=0; i<obj->group->done_op_count; i++)
+ for(i=0; i<(obj->group->done_op_count-1); i++)
{
- aesop_sem_up(&obj->group->sem);
+ ret = aesop_sem_up(&obj->group->sem);
+ assert(ret == AE_SUCCESS);
}
}
diff --git a/code/tests/Makefile.subdir b/code/tests/Makefile.subdir
index 8912e05..e95af09 100644
--- a/code/tests/Makefile.subdir
+++ b/code/tests/Makefile.subdir
@@ -18,6 +18,7 @@ TESTS += \
tests/triton-noop-multi-svr.sh \
tests/triton-touch.sh \
tests/triton-cp.sh \
+ tests/triton-cp-big.sh \
tests/triton-ls.sh \
tests/triton-rm.sh \
tests/triton-show-system-state.sh \
diff --git a/code/tests/triton-cp-big.sh b/code/tests/triton-cp-big.sh
new file mode 100755
index 0000000..ba46667
--- /dev/null
+++ b/code/tests/triton-cp-big.sh
@@ -0,0 +1,70 @@
+#!/bin/bash
+
+if [ -z $srcdir ]; then
+ echo srcdir variable not set.
+ exit 1
+fi
+source $srcdir/tests/test-util.sh
+
+# start 4 servers with 15 second wait, 120s timeout
+test_start_servers 4 15 120
+
+# actual test case
+#####################
+
+# generate data file. We start with 1000000 bytes of random data and
+# then create a repeating pattern. The original random data is aligned on
+# power of 10 rather than power of 2 boundary on purpose to make sure that
+# we can detect errors that might occur on pipeline buffer boundaries.
+#
+dd if=/dev/urandom of=/tmp/128-$$.dat bs=1000000 count=1
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=1 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=2 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=4 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=8 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=16 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=32 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=64 conv=notrunc oflag=append
+
+run_to 90 src/admin-tools/triton-cp $svr1 posix:/tmp/128-$$.dat ::128 3
+if [ $? -ne 0 ]; then
+ rm -rf /tmp/128-$$.dat
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+run_to 90 src/admin-tools/triton-cp $svr1 ::128 posix:/tmp/128-$$.dat.check
+if [ $? -ne 0 ]; then
+ rm -rf /tmp/128-$$.dat /tmp/128-$$.dat.check
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+run_to 90 src/admin-tools/triton-rm --server $svr1 ::128
+if [ $? -ne 0 ]; then
+ rm -rf /tmp/128-$$.dat /tmp/128-$$.dat.check
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+diff /tmp/128-$$.dat /tmp/128-$$.dat.check
+if [ $? -ne 0 ]; then
+ echo ERROR: data copied in and out of Triton does not match.
+ echo ERROR: see files /tmp/128-$$.dat and /tmp/128-$$.dat.check
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+rm -rf /tmp/128-$$.dat /tmp/128-$$.dat.check
+
+#####################
+
+# tear down
+run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+
+wait
+exit 0
diff --git a/code/tests/triton-cp.sh b/code/tests/triton-cp.sh
index 099ff50..39f9607 100755
--- a/code/tests/triton-cp.sh
+++ b/code/tests/triton-cp.sh
@@ -13,22 +13,25 @@ test_start_servers 4 15 60
#####################
dd if=/dev/urandom of=/tmp/8-$$.dat bs=1K count=8
-run_to 60 src/admin-tools/triton-cp $svr1 posix:/tmp/8-$$.dat 1.8 3
+run_to 60 src/admin-tools/triton-cp $svr1 posix:/tmp/8-$$.dat ::8 3
if [ $? -ne 0 ]; then
+ rm -rf /tmp/8-$$.dat
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
fi
-run_to 60 src/admin-tools/triton-cp $svr1 1.8 posix:/tmp/8-$$.dat.check
+run_to 60 src/admin-tools/triton-cp $svr1 ::8 posix:/tmp/8-$$.dat.check
if [ $? -ne 0 ]; then
+ rm -rf /tmp/8-$$.dat /tmp/8-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
fi
-run_to 60 src/admin-tools/triton-rm --server $svr1 1.8
+run_to 60 src/admin-tools/triton-rm --server $svr1 ::8
if [ $? -ne 0 ]; then
+ rm -rf /tmp/8-$$.dat /tmp/8-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
@@ -37,6 +40,7 @@ fi
diff /tmp/8-$$.dat /tmp/8-$$.dat.check
if [ $? -ne 0 ]; then
echo ERROR: data copied in and out of Triton does not match.
+ echo ERROR: see files /tmp/128-$$.dat and /tmp/128-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
@@ -45,6 +49,7 @@ fi
#####################
# tear down
+rm -rf /tmp/8-$$.dat /tmp/8-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. 1f4b68d19d951a4732a32eea456b22f28f31c4f1
by noreply@mcs.anl.gov 27 Apr '14
by noreply@mcs.anl.gov 27 Apr '14
27 Apr '14
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, trac-277-bulkio has been updated
via 1f4b68d19d951a4732a32eea456b22f28f31c4f1 (commit)
via 948907e0c766b0f5e1f4bad1ea9de33bf2de19b7 (commit)
via f59a367b2b8334fb28318b42ffe3ee8e483ce8b4 (commit)
via 7e355cb8968e2e785867e702ee9073d23ea14a32 (commit)
from bbdbb1229d4281409040666ba0489486d809d6d0 (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 1f4b68d19d951a4732a32eea456b22f28f31c4f1
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 21:34:54 2014 -0400
clean up 64 bit integer usage
commit 948907e0c766b0f5e1f4bad1ea9de33bf2de19b7
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 21:05:35 2014 -0400
comment out some of the most excessive dbg msgs
commit f59a367b2b8334fb28318b42ffe3ee8e483ce8b4
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 21:04:03 2014 -0400
error handling in rosd-write
commit 7e355cb8968e2e785867e702ee9073d23ea14a32
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 20:56:41 2014 -0400
error handling in rosd-read
-----------------------------------------------------------------------
Summary of changes:
code/src/replicated-osd/buffer-mgmt.ae | 6 +-
code/src/replicated-osd/buffer-mgmt.hae | 2 +-
code/src/replicated-osd/rosd-read.ae | 49 +++++++++++++-------
code/src/replicated-osd/rosd-write.ae | 75 +++++++++++++++++++-----------
code/src/replicated-osd/rosd.ae | 8 ++--
code/src/replicated-osd/rosd.hae | 4 +-
6 files changed, 90 insertions(+), 54 deletions(-)
Diff of changes:
diff --git a/code/src/replicated-osd/buffer-mgmt.ae b/code/src/replicated-osd/buffer-mgmt.ae
index c5c8e91..c75ca09 100644
--- a/code/src/replicated-osd/buffer-mgmt.ae
+++ b/code/src/replicated-osd/buffer-mgmt.ae
@@ -15,7 +15,7 @@ struct buffer_mgmt_instance
struct buffer_mgmt_token
{
char* buffer;
- int size;
+ int64_t size;
struct buffer_mgmt_instance* instance;
};
@@ -72,10 +72,10 @@ struct buffer_mgmt_instance* buffer_mgmt_create_instance(
__blocking triton_ret_t buffer_mgmt_alloc(
struct buffer_mgmt_instance *instance,
- int requested_size, int* allocated_size,
+ int64_t requested_size, int64_t* allocated_size,
char** buffer, struct buffer_mgmt_token** token)
{
- int size_to_alloc = 0;
+ int64_t size_to_alloc = 0;
int ret;
assert(requested_size > 0);
diff --git a/code/src/replicated-osd/buffer-mgmt.hae b/code/src/replicated-osd/buffer-mgmt.hae
index 88fbc2f..d887ab1 100644
--- a/code/src/replicated-osd/buffer-mgmt.hae
+++ b/code/src/replicated-osd/buffer-mgmt.hae
@@ -14,7 +14,7 @@ struct buffer_mgmt_instance* buffer_mgmt_create_instance(
__blocking triton_ret_t buffer_mgmt_alloc(
struct buffer_mgmt_instance* instance,
- int requested_size, int* allocated_size,
+ int64_t requested_size, int64_t* allocated_size,
char** buffer, struct buffer_mgmt_token** token);
void buffer_mgmt_free(struct buffer_mgmt_token* token);
diff --git a/code/src/replicated-osd/rosd-read.ae b/code/src/replicated-osd/rosd-read.ae
index 022b033..e0f05c7 100644
--- a/code/src/replicated-osd/rosd-read.ae
+++ b/code/src/replicated-osd/rosd-read.ae
@@ -200,8 +200,7 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
int my_position;
int ret;
na_addr_t src_addr;
- /* TODO: 64 bit? */
- int size_remaining = 0;
+ int64_t size_remaining = 0;
int64_t remote_offset = 0;
int64_t local_offset = 0;
aesop_sem_t sem;
@@ -244,8 +243,7 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
pwait
{
pprivate int i = 0;
- /* TODO: 64 bit? */
- pprivate int this_size = 0;
+ pprivate int64_t this_size = 0;
pprivate int64_t this_out_size = 0;
pprivate int64_t this_local_offset = 0;
pprivate int64_t this_remote_offset = 0;
@@ -277,8 +275,18 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
&this_remote_offset,
&token,
&buffer);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
+ triton_mutex_lock(&output_mutex);
+ if(triton_is_error(tret) && !triton_is_error(out.tret))
+ {
+ out.tret = tret;
+ this_size = 0;
+ }
+ else if(triton_is_error(tret))
+ {
+ triton_error_destroy(tret);
+ }
+ triton_mutex_unlock(&output_mutex);
+
//triton_debug(triton_dbg_rosd, "triton_rpc_rosd_read() pbranch %d finished setup_one_buffer(), size_remaining: %d.\n", i, size_remaining);
if(this_size > 0)
@@ -296,21 +304,30 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
in.bulk_handle,
&this_out_size,
&this_txn_number);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
//triton_debug(triton_dbg_rosd, "triton_rpc_rosd_read() pbranch %d finished xfer_one_buffer().\n", i);
/* accumulate results in rpc response */
triton_mutex_lock(&output_mutex);
- out.out_size += this_out_size;
- if(out.txn_number == 0 && out.txn_number == this_txn_number)
- out.txn_number = this_txn_number;
+ if(triton_is_error(tret) && !triton_is_error(out.tret))
+ {
+ out.tret = tret;
+ }
+ else if(triton_is_error(tret))
+ {
+ triton_error_destroy(tret);
+ }
else
{
- /* TODO: what are we supposed to do on mixed txn
- * numbers?
- */
- out.txn_number = 0;
+ out.out_size += this_out_size;
+ if(out.txn_number == 0 && out.txn_number == this_txn_number)
+ out.txn_number = this_txn_number;
+ else
+ {
+ /* TODO: what are we supposed to do on mixed txn
+ * numbers?
+ */
+ out.txn_number = 0;
+ }
}
triton_mutex_unlock(&output_mutex);
}
@@ -318,7 +335,7 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
{
this_out_size = this_size;
}
- }while(this_out_size > 0);
+ }while(this_out_size > 0 && !triton_is_error(out.tret));
if(token != NULL)
{
diff --git a/code/src/replicated-osd/rosd-write.ae b/code/src/replicated-osd/rosd-write.ae
index 5d89e95..776788d 100644
--- a/code/src/replicated-osd/rosd-write.ae
+++ b/code/src/replicated-osd/rosd-write.ae
@@ -151,8 +151,8 @@ static __blocking triton_ret_t rosd_write_run_iteration(
int from_client_flag,
triton_rpc_rosd_write_in_t* in,
aesop_sem_t *sem,
- int* size_remaining,
- int* this_size,
+ int64_t* size_remaining,
+ int64_t* this_size,
int64_t *remote_offset,
int64_t *local_offset,
struct buffer_mgmt_token **token,
@@ -170,7 +170,7 @@ static __blocking triton_ret_t rosd_write_run_iteration(
/* note that this function protects shared variables
* using a semaphore
*/
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling setup_one_buffer().\n");
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling setup_one_buffer().\n");
tret = rosd_pipeline_setup_one_buffer(
rosd_write_buffers,
sem,
@@ -182,14 +182,16 @@ static __blocking triton_ret_t rosd_write_run_iteration(
&this_remote_offset,
token,
buffer);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished setup_one_buffer(), this_size: %d, this_local_offset: %ld, this_remote_offset %ld.\n", *this_size, this_local_offset, this_remote_offset);
+ if(triton_is_error(tret))
+ {
+ return(tret);
+ }
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished setup_one_buffer(), this_size: %d, this_local_offset: %ld, this_remote_offset %ld.\n", *this_size, this_local_offset, this_remote_offset);
if(*this_size > 0)
{
/* perform buffer transfer */
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling pull_one_buffer().\n");
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling pull_one_buffer().\n");
tret = rosd_write_pull_one_buffer(
src_addr,
*buffer,
@@ -198,18 +200,23 @@ static __blocking triton_ret_t rosd_write_run_iteration(
*this_size,
this_remote_offset,
in->bulk_handle);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished pull_one_buffer().\n");
+ if(triton_is_error(tret))
+ {
+ return(tret);
+ }
+
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished pull_one_buffer().\n");
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch starting do_work().\n");
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch starting do_work().\n");
tret = rosd_write_do_work(next_addr, in->oid,
in->oid_fork, *buffer, *this_size, this_local_offset,
in->flags, in->replication_factor, in->txn_number,
my_position, from_client_flag);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished do_work().\n");
+ if(triton_is_error(tret))
+ {
+ return(tret);
+ }
+ //triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished do_work().\n");
*more_flag = 1;
}
@@ -225,17 +232,20 @@ static __blocking triton_ret_t rosd_write_run_pipeline(
triton_rpc_rosd_write_in_t* in
)
{
- /* TODO: 64 bit? */
- int size_remaining = 0;
+ int64_t size_remaining = 0;
int64_t remote_offset = 0;
int64_t local_offset = 0;
aesop_sem_t sem;
+ triton_ret_t out_tret = TRITON_SUCCESS;
+ triton_mutex_t output_mutex;
size_remaining = in->size;
local_offset = in->offset;
remote_offset = 0;
aesop_sem_init(&sem, 1);
+ triton_mutex_init(&output_mutex, NULL);
+
pwait
{
pprivate int i = 0;
@@ -243,17 +253,17 @@ static __blocking triton_ret_t rosd_write_run_pipeline(
pprivate struct buffer_mgmt_token *token = NULL;
pprivate char* buffer = NULL;
pprivate int more_flag = 1;
- pprivate int this_size = 0;
+ pprivate int64_t this_size = 0;
- triton_debug(triton_dbg_rosd,
- "rosd_write_run_pipeline() starting pwait, size_remaining: %d.\n",
- size_remaining);
+ //triton_debug(triton_dbg_rosd,
+ // "rosd_write_run_pipeline() starting pwait, size_remaining: %d.\n",
+ // size_remaining);
for(i=0; i<ROSD_WRITE_XFER_PIPELINE_DEPTH; i++)
{
pbranch
{
- triton_debug(triton_dbg_rosd,
- "rosd_write_run_pipeline() starting pbranch %d.\n", i);
+ //triton_debug(triton_dbg_rosd,
+ // "rosd_write_run_pipeline() starting pbranch %d.\n", i);
do
{
@@ -261,22 +271,31 @@ static __blocking triton_ret_t rosd_write_run_pipeline(
my_position, from_client_flag, in, &sem,
&size_remaining, &this_size, &remote_offset,
&local_offset, &token, &buffer, &more_flag);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
- }while(more_flag);
+ }while(more_flag && !triton_is_error(tret));
if(token != NULL)
{
buffer_mgmt_free(token);
}
- triton_debug(triton_dbg_rosd,
- "rosd_write_run_pipeline() finishing pbranch %d.\n", i);
+ //triton_debug(triton_dbg_rosd,
+ // "rosd_write_run_pipeline() finishing pbranch %d.\n", i);
+ triton_mutex_lock(&output_mutex);
+ if(triton_is_error(tret) && !triton_is_error(out_tret))
+ {
+ out_tret = tret;
+ }
+ else if(triton_is_error(tret))
+ {
+ triton_error_destroy(tret);
+ }
+ triton_mutex_unlock(&output_mutex);
}
}
}
aesop_sem_destroy(&sem);
+ triton_mutex_destroy(&output_mutex);
- return(TRITON_SUCCESS);
+ return(out_tret);
}
diff --git a/code/src/replicated-osd/rosd.ae b/code/src/replicated-osd/rosd.ae
index f70c89e..30424b6 100644
--- a/code/src/replicated-osd/rosd.ae
+++ b/code/src/replicated-osd/rosd.ae
@@ -388,10 +388,10 @@ __blocking void trigger_server_fault(triton_ret_t tret)
__blocking triton_ret_t rosd_pipeline_setup_one_buffer(
struct buffer_mgmt_instance* buffer_instance,
aesop_sem_t *sem,
- int* size_remaining,
+ int64_t* size_remaining,
int64_t *remote_offset,
int64_t *local_offset,
- int* this_size,
+ int64_t* this_size,
int64_t* this_local_offset,
int64_t* this_remote_offset,
struct buffer_mgmt_token **token,
@@ -410,7 +410,7 @@ __blocking triton_ret_t rosd_pipeline_setup_one_buffer(
return(TRITON_ERR_UNKNOWN);
}
- triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() acquired semaphore with size_remaining %d\n", *size_remaining);
+ triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() acquired semaphore with size_remaining %lld\n", lld(*size_remaining));
if((*size_remaining) == 0)
{
@@ -454,7 +454,7 @@ __blocking triton_ret_t rosd_pipeline_setup_one_buffer(
*local_offset += *this_size;
*remote_offset += *this_size;
- triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() releasing semaphore with size_remaining %d\n", *size_remaining);
+ triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() releasing semaphore with size_remaining %lld\n", lld(*size_remaining));
ret = aesop_sem_up(sem);
assert(ret == AE_SUCCESS);
diff --git a/code/src/replicated-osd/rosd.hae b/code/src/replicated-osd/rosd.hae
index da89f27..633fea8 100644
--- a/code/src/replicated-osd/rosd.hae
+++ b/code/src/replicated-osd/rosd.hae
@@ -67,10 +67,10 @@ __blocking triton_ret_t remote_triton_rpc_rosd_read(
__blocking triton_ret_t rosd_pipeline_setup_one_buffer(
struct buffer_mgmt_instance* buffer_instance,
aesop_sem_t *sem,
- int* size_remaining,
+ int64_t* size_remaining,
int64_t *remote_offset,
int64_t *local_offset,
- int* this_size,
+ int64_t* this_size,
int64_t* this_local_offset,
int64_t* this_remote_offset,
struct buffer_mgmt_token **token,
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. bbdbb1229d4281409040666ba0489486d809d6d0
by noreply@mcs.anl.gov 27 Apr '14
by noreply@mcs.anl.gov 27 Apr '14
27 Apr '14
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, trac-277-bulkio has been updated
via bbdbb1229d4281409040666ba0489486d809d6d0 (commit)
via d0ecd3fbaf8c5fe3df48672b7482016870f1df0d (commit)
via 8df78189eec3cce18fea516609f409e4345812dd (commit)
via 8860feb91cafe6f6be8ee166c33d4eec7da62503 (commit)
from 7862829659b1f8844f2e058aad99d85f38c2bfcf (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 bbdbb1229d4281409040666ba0489486d809d6d0
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 20:48:30 2014 -0400
clean up temporary files in triton-cp.sh
commit d0ecd3fbaf8c5fe3df48672b7482016870f1df0d
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 20:43:08 2014 -0400
fix oid format in triton-cp.sh
commit 8df78189eec3cce18fea516609f409e4345812dd
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 20:39:21 2014 -0400
clean up temporary filesin triton-cp-big.sh
commit 8860feb91cafe6f6be8ee166c33d4eec7da62503
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Sun Apr 27 20:36:25 2014 -0400
fix minor tosd logic bug
- semaphore up() too many times when sync coalescing
-----------------------------------------------------------------------
Summary of changes:
code/src/transactional-osd/transactional-osd.ae | 5 +++--
code/tests/triton-cp-big.sh | 6 ++++++
code/tests/triton-cp.sh | 11 ++++++++---
3 files changed, 17 insertions(+), 5 deletions(-)
Diff of changes:
diff --git a/code/src/transactional-osd/transactional-osd.ae b/code/src/transactional-osd/transactional-osd.ae
index 0d7b900..65aa8c8 100644
--- a/code/src/transactional-osd/transactional-osd.ae
+++ b/code/src/transactional-osd/transactional-osd.ae
@@ -424,9 +424,10 @@ static __blocking triton_ret_t end_write(struct coalesce_obj *obj)
}
/* notify everyone else in group */
- for(i=0; i<obj->group->done_op_count; i++)
+ for(i=0; i<(obj->group->done_op_count-1); i++)
{
- aesop_sem_up(&obj->group->sem);
+ ret = aesop_sem_up(&obj->group->sem);
+ assert(ret == AE_SUCCESS);
}
}
diff --git a/code/tests/triton-cp-big.sh b/code/tests/triton-cp-big.sh
index 60a76d7..ba46667 100755
--- a/code/tests/triton-cp-big.sh
+++ b/code/tests/triton-cp-big.sh
@@ -28,6 +28,7 @@ dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=64 conv=notrunc oflag=appen
run_to 90 src/admin-tools/triton-cp $svr1 posix:/tmp/128-$$.dat ::128 3
if [ $? -ne 0 ]; then
+ rm -rf /tmp/128-$$.dat
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
@@ -35,6 +36,7 @@ fi
run_to 90 src/admin-tools/triton-cp $svr1 ::128 posix:/tmp/128-$$.dat.check
if [ $? -ne 0 ]; then
+ rm -rf /tmp/128-$$.dat /tmp/128-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
@@ -42,6 +44,7 @@ fi
run_to 90 src/admin-tools/triton-rm --server $svr1 ::128
if [ $? -ne 0 ]; then
+ rm -rf /tmp/128-$$.dat /tmp/128-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
@@ -50,11 +53,14 @@ fi
diff /tmp/128-$$.dat /tmp/128-$$.dat.check
if [ $? -ne 0 ]; then
echo ERROR: data copied in and out of Triton does not match.
+ echo ERROR: see files /tmp/128-$$.dat and /tmp/128-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
fi
+rm -rf /tmp/128-$$.dat /tmp/128-$$.dat.check
+
#####################
# tear down
diff --git a/code/tests/triton-cp.sh b/code/tests/triton-cp.sh
index 099ff50..39f9607 100755
--- a/code/tests/triton-cp.sh
+++ b/code/tests/triton-cp.sh
@@ -13,22 +13,25 @@ test_start_servers 4 15 60
#####################
dd if=/dev/urandom of=/tmp/8-$$.dat bs=1K count=8
-run_to 60 src/admin-tools/triton-cp $svr1 posix:/tmp/8-$$.dat 1.8 3
+run_to 60 src/admin-tools/triton-cp $svr1 posix:/tmp/8-$$.dat ::8 3
if [ $? -ne 0 ]; then
+ rm -rf /tmp/8-$$.dat
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
fi
-run_to 60 src/admin-tools/triton-cp $svr1 1.8 posix:/tmp/8-$$.dat.check
+run_to 60 src/admin-tools/triton-cp $svr1 ::8 posix:/tmp/8-$$.dat.check
if [ $? -ne 0 ]; then
+ rm -rf /tmp/8-$$.dat /tmp/8-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
fi
-run_to 60 src/admin-tools/triton-rm --server $svr1 1.8
+run_to 60 src/admin-tools/triton-rm --server $svr1 ::8
if [ $? -ne 0 ]; then
+ rm -rf /tmp/8-$$.dat /tmp/8-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
@@ -37,6 +40,7 @@ fi
diff /tmp/8-$$.dat /tmp/8-$$.dat.check
if [ $? -ne 0 ]; then
echo ERROR: data copied in and out of Triton does not match.
+ echo ERROR: see files /tmp/128-$$.dat and /tmp/128-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
exit 1
@@ -45,6 +49,7 @@ fi
#####################
# tear down
+rm -rf /tmp/8-$$.dat /tmp/8-$$.dat.check
run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
wait
hooks/post-receive
--
1
0
25 Apr '14
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 c347c633e3ae197603ebd7ae3428935efb2d8491 (commit)
via 24e05f0d389bc131a1067fd24e99e16eef8e67a4 (commit)
from ed264ff8dfe6584f4c0c6d81bbb9a2cc8cf91456 (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 c347c633e3ae197603ebd7ae3428935efb2d8491
Merge: 24e05f0d389bc131a1067fd24e99e16eef8e67a4 ed264ff8dfe6584f4c0c6d81bbb9a2cc8cf91456
Author: Kevin Harms <harms(a)alcf.anl.gov>
Date: Fri Apr 25 17:32:48 2014 -0500
Merge branch 'master' of git.mcs.anl.gov:triton
commit 24e05f0d389bc131a1067fd24e99e16eef8e67a4
Author: Kevin Harms <harms(a)alcf.anl.gov>
Date: Fri Apr 25 17:31:30 2014 -0500
Update to latest ASGheader. Add C wrapper libraries to asg aesop interfaces
-----------------------------------------------------------------------
Summary of changes:
code/include/asg.h | 29 +++++++---
code/src/asg/Makefile.subdir | 4 +-
code/src/asg/asg-internal.ae | 51 +++++++++--------
code/src/asg/asg-internal.hae | 77 ++-------------------------
code/src/asg/asg.c | 118 ++++++++++++++++++++++++++++++++++++++---
5 files changed, 165 insertions(+), 114 deletions(-)
Diff of changes:
diff --git a/code/include/asg.h b/code/include/asg.h
index b74263a..8544925 100644
--- a/code/include/asg.h
+++ b/code/include/asg.h
@@ -4,6 +4,11 @@
#include <stdint.h>
#include <string.h>
+#ifdef __cplusplus
+extern "C"
+{
+#endif
+
typedef intptr_t asg_instance_t;
typedef intptr_t asg_location_t;
@@ -40,6 +45,7 @@ typedef uint64_t asg_version_t;
typedef enum
{
ASG_SUCCESS = 0,
+ ASG_ERR_VERSION = 1,
ASG_ERR_OTHER
} asg_ret_t;
@@ -53,7 +59,7 @@ typedef enum
/** For reading: read in order until a record is found which has the
* same or higher version number than the one specified.
* For writing/punch: update records in order as specified until the end of
- * the specified range is encounder or until a record is found which has a
+ * the specified range is encountered or until a record is found which has a
* version greater or equal to the one specified. */
ASG_COND_UNTIL = 0x0001,
@@ -64,11 +70,16 @@ typedef enum
*/
ASG_COND_ALL = 0x0002,
+ /** Always read/write */
+ ASG_COND_NONE = 0x0000,
+ ASG_COND_UNCONDITIONAL = ASG_COND_NONE,
/** For writing/punch: automatically set the version of the updated records
* to a version number guaranteed to be higher than any of the previously
* existing version numbers in the range. */
- ASG_AUTO_VERSION = 0x0004
+ ASG_AUTO_VERSION = 0x0004,
+
+
} asg_flags_t;
@@ -79,7 +90,7 @@ int asg_initialize (asg_instance_t * instance, const char * options);
/** Close storage system instance */
-int asg_finalize (void);
+int asg_finalize (asg_instance_t instance);
/**
* Retrieve data
@@ -114,12 +125,12 @@ int asg_write (
asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
const void * data,
asg_size_t * transferred);
/**
- * Punched is exactly like write but writes zero length records.
+ * Punch is exactly like write but writes zero length records.
*
* Transferred indicates the number of records updated.
*/
@@ -131,10 +142,9 @@ int asg_punch (
asg_fork_id_t fork,
asg_record_id_t start_record,
asg_size_t recordcount,
- asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
asg_size_t * transferred);
@@ -183,7 +193,7 @@ int asg_reset (
* which information can be retrieved (and this value can be passed to 'start'
* in subsequent calls).
*
- * If no more items are availalbe, *next is set to ASG_xxx_NULL
+ * If no more items are available, *next is set to ASG_xxx_NULL
*
* Probe does not create a snapshot.
*/
@@ -262,5 +272,8 @@ int asg_probe_fork (
asg_size_t * transferred,
asg_record_id_t * next);
+#ifdef __cplusplus
+}
+#endif
#endif
diff --git a/code/src/asg/Makefile.subdir b/code/src/asg/Makefile.subdir
index 225aa89..4e5e677 100644
--- a/code/src/asg/Makefile.subdir
+++ b/code/src/asg/Makefile.subdir
@@ -1,7 +1,7 @@
src_libtriton_a_SOURCES += \
src/asg/asg.c \
- src/asg/asg-internal.ae\
- src/asg/asg-internal.hae
+ src/asg/asg-internal.ae \
+ src/asg/asg-internal.h
AE_SRC += \
src/asg/asg-internal.ae
diff --git a/code/src/asg/asg-internal.ae b/code/src/asg/asg-internal.ae
index 5e6653d..fa7529b 100644
--- a/code/src/asg/asg-internal.ae
+++ b/code/src/asg/asg-internal.ae
@@ -15,10 +15,10 @@
#include "src/common/triton-bootstrap.hae"
#include "src/system-state/system-state.hae"
#include "src/replicated-osd/rosd.hae"
+#include "src/asg/asg-internal.hae"
#include "asg.h"
-static int replication_factor = 1;
static triton_node_t *node_array = NULL;
static int node_array_count;
@@ -26,13 +26,13 @@ static int node_array_count;
* should initialize all the server name and their address
* first step, use options as the server address we want to connect
*/
-__blocking triton_ret_t asg_i_initialize(asg_instance_t *instance,
+__blocking int asg_i_initialize(asg_instance_t *instance,
const char *options)
{
triton_ret_t tret;
//if initialized before, return;
if (node_array != NULL)
- return(TRITON_SUCCESS);
+ return(ASG_SUCCESS);
//initialize error to /dev/stderr
tret = triton_debug_init();
@@ -40,7 +40,7 @@ __blocking triton_ret_t asg_i_initialize(asg_instance_t *instance,
{
triton_error_print(tret, "triton_debug_init()");
triton_error_destroy(tret);
- return(-1);
+ return(ASG_ERR_OTHER);
}
triton_debug_enable("/dev/stderr", "none");
@@ -50,7 +50,7 @@ __blocking triton_ret_t asg_i_initialize(asg_instance_t *instance,
{
triton_error_print(tret, "triton_mercury_engine_init()");
triton_error_destroy(tret);
- return(-1);
+ return(ASG_ERR_OTHER);
}
triton_rpc_system_state_register();
@@ -61,16 +61,17 @@ __blocking triton_ret_t asg_i_initialize(asg_instance_t *instance,
{
triton_error_print(tret, "triton_attach()");
triton_error_destroy(tret);
- return(-1);
+ return(ASG_ERR_OTHER);
}
tret = system_state_list_nodes("status", &node_array, &node_array_count);
if (triton_is_error(tret))
{
- return(tret);
+ triton_error_destroy(tret);
+ return(ASG_ERR_OTHER);
}
- return(TRITON_SUCCESS);
+ return(ASG_SUCCESS);
}
/**
* dismiss location arg.
@@ -78,7 +79,7 @@ __blocking triton_ret_t asg_i_initialize(asg_instance_t *instance,
* fork => oid_fork;
* record => bytes stream
*/
-__blocking triton_ret_t asg_i_write (
+__blocking int asg_i_write (
asg_instance_t instance,
asg_location_t location,
asg_container_id_t container,
@@ -89,7 +90,7 @@ __blocking triton_ret_t asg_i_write (
asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
const void * data,
asg_size_t * transferred)
{
@@ -117,7 +118,7 @@ __blocking triton_ret_t asg_i_write (
{
triton_error_print(tret, "remote_triton_rpc_rosd_create");
triton_error_destroy(tret);
- return(-1);
+ return(ASG_ERR_OTHER);
}
tret = remote_triton_rpc_rosd_write(oid,
@@ -132,10 +133,10 @@ __blocking triton_ret_t asg_i_write (
{
triton_error_print(tret, "remote_triton_rpc_rosd_write ");
triton_error_destroy(tret);
- return(-1);
+ return(ASG_ERR_OTHER);
}
triton_error_destroy(tret);
- return(TRITON_SUCCESS);
+ return(ASG_SUCCESS);
}
__blocking int asg_i_read (
@@ -155,10 +156,10 @@ __blocking int asg_i_read (
{
triton_ret_t tret;
- //build oid and get server info from requests;
uint128_t oid;
- oid.u = container;
- oid.l = object;
+
+ oid.u = object;
+ oid.l = container;
tret = remote_triton_rpc_rosd_read(
oid,
@@ -167,17 +168,19 @@ __blocking int asg_i_read (
bufsize,
start_record,
version_info,
- transferred,
+ (int64_t*)transferred,
flags);
+
if (triton_is_error(tret))
{
triton_error_print(tret, "remote_triton_rpc_rosd_read ");
triton_error_destroy(tret);
- return(-1);
+ return(ASG_ERR_OTHER);
}
+
triton_error_destroy(tret);
- return(TRITON_SUCCESS);
+ return(ASG_SUCCESS);
}
__blocking int asg_i_probe_fork(
@@ -209,7 +212,7 @@ __blocking int asg_i_punch (
asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
asg_size_t * transferred)
{
triton_ret_t tret;
@@ -232,14 +235,14 @@ __blocking int asg_i_punch (
{
triton_error_print(tret, "remote_triton_rpc_rosd_write ");
triton_error_destroy(tret);
- return(-1);
+ return(ASG_ERR_OTHER);
}
triton_error_destroy(tret);
- return(TRITON_SUCCESS);
+ return(ASG_SUCCESS);
}
-__blocking int asg_i_finalize(void)
+__blocking int asg_i_finalize(asg_instance_t instance)
{
triton_mercury_engine_finalize();
- return(TRITON_SUCCESS);
+ return(ASG_SUCCESS);
}
diff --git a/code/src/asg/asg-internal.hae b/code/src/asg/asg-internal.hae
index 60c3e8a..1a639db 100644
--- a/code/src/asg/asg-internal.hae
+++ b/code/src/asg/asg-internal.hae
@@ -8,79 +8,10 @@
#define __ASG_INTERNAL_HAE__
-#include <stdint.h>
-#include <string.h>
-
-typedef intptr_t asg_instance_t;
-
-typedef intptr_t asg_location_t;
-
-typedef uint64_t asg_container_id_t;
-typedef uint64_t asg_object_id_t;
-typedef uint64_t asg_fork_id_t;
-typedef uint64_t asg_record_id_t;
-
-typedef uint64_t asg_offset_t;
-typedef uint64_t asg_size_t;
-
-
-typedef uint64_t asg_version_t;
-
-/** The highest valid version number */
-#define ASG_VERSION_MAX (2^64-2)
-
-/** Constant indicating multiple versions
- * (used in reading) */
-#define ASG_VERSION_MIXED (2^64-1)
-
-
-#define ASG_LOCATION_AUTO 0
-
-#define ASG_RECORD_NULL 0
-#define ASG_FORK_NULL 0
-#define ASG_OBJECT_NULL 0
-#define ASG_CONTAINER_NULL 0
-
-/**
- * Return/error codes
- */
-typedef enum
-{
- ASG_SUCCESS = 0,
- ASG_ERR_OTHER
-} asg_ret_t;
-
-
-
-/**
- * Flags for use in read,write,punch.
- */
-typedef enum
-{
- /** For reading: read in order until a record is found which has the
- * same or higher version number than the one specified.
- * For writing/punch: update records in order as specified until the end of
- * the specified range is encounder or until a record is found which has a
- * version greater or equal to the one specified. */
- ASG_COND_UNTIL = 0x0001,
-
- /** For reading: only read if *all* records in the specified range
- * have a version number strictly smaller than the one specified.
- * For writing/punch: only write if *all* records in the specified range
- * have a version number strictly smaller than the one specified.
- */
- ASG_COND_ALL = 0x0002,
-
-
- /** For writing/punch: automatically set the version of the updated records
- * to a version number guaranteed to be higher than any of the previously
- * existing version numbers in the range. */
- ASG_AUTO_VERSION = 0x0004
-
-} asg_flags_t;
+#include <asg.h>
__blocking int asg_i_initialize(asg_instance_t *instance, const char *options);
-__blocking int asg_i_finalize(void);
+__blocking int asg_i_finalize(asg_instance_t instance);
__blocking int asg_i_write (
asg_instance_t instance,
@@ -93,7 +24,7 @@ __blocking int asg_i_write (
asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
const void * data,
asg_size_t * transferred);
@@ -128,7 +59,7 @@ __blocking int asg_i_punch (
asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
asg_size_t * transferred);
#endif
diff --git a/code/src/asg/asg.c b/code/src/asg/asg.c
index 558bae9..52bc8a2 100644
--- a/code/src/asg/asg.c
+++ b/code/src/asg/asg.c
@@ -1,12 +1,72 @@
-#include "asg.h"
+#include <aesop/aesop.h>
+#include <src/common/triton-error.h>
+#include <asg.h>
+#include <src/asg/asg-internal.h>
+
+#define AE_POLL_TIMEOUT 1000
+
+typedef struct asg_completion_s
+{
+ int ret;
+ int done;
+} asg_completion_t;
+
+typedef struct asg_op_id_s
+{
+ asg_completion_t c;
+ ae_hints_t hints;
+} asg_op_id_t;
+
+static void asg_callback (void *user_ptr, int ret)
+{
+ asg_completion_t *c = (asg_completion_t *) user_ptr;
+ c->ret = ret;
+ c->done = 1;
+ ae_poll_break();
+ return;
+}
int asg_initialize (asg_instance_t * instance, const char * options)
{
- return ASG_SUCCESS;
+ asg_completion_t c;
+ ae_hints_t hints;
+ ae_op_id_t op_id;
+ int r;
+
+ c.done = 0;
+ ae_hints_init(&hints);
+
+ r = ext_post_blocking(asg_i_initialize,
+ asg_callback,
+ &c,
+ &hints,
+ &op_id,
+ &c.ret,
+ instance,
+ options);
+ if (r == AE_SUCCESS)
+ {
+ while(!c.done)
+ {
+ ae_poll(AE_POLL_TIMEOUT);
+ }
+ }
+ else if (r == AE_IMMEDIATE_COMPLETION)
+ {
+ // nothing to do here
+ }
+ else
+ {
+ c.ret = ASG_ERR_OTHER;
+ }
+
+ ae_hints_destroy(&hints);
+
+ return c.ret;
}
-int asg_finalize (void)
+int asg_finalize (asg_instance_t instance)
{
return ASG_SUCCESS;
}
@@ -26,7 +86,52 @@ int asg_read (
asg_size_t * transferred,
asg_version_t * version_info)
{
- return ASG_SUCCESS;
+ asg_completion_t c;
+ ae_hints_t hints;
+ ae_op_id_t op_id;
+ int r;
+
+ c.done = 0;
+ ae_hints_init(&hints);
+
+ r = ext_post_blocking(asg_i_read,
+ asg_callback,
+ &c,
+ &hints,
+ &op_id,
+ &c.ret,
+ instance,
+ location,
+ container,
+ object,
+ fork,
+ start_record,
+ recordcount,
+ flags,
+ version_condition,
+ buf,
+ bufsize,
+ transferred,
+ version_info);
+ if (r == AE_SUCCESS)
+ {
+ while(!c.done)
+ {
+ ae_poll(AE_POLL_TIMEOUT);
+ }
+ }
+ else if (r == AE_IMMEDIATE_COMPLETION)
+ {
+ // nothing to do here
+ }
+ else
+ {
+ c.ret = ASG_ERR_OTHER;
+ }
+
+ ae_hints_destroy(&hints);
+
+ return c.ret;
}
int asg_write (
@@ -40,7 +145,7 @@ int asg_write (
asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
const void * data,
asg_size_t * transferred)
{
@@ -55,10 +160,9 @@ int asg_punch (
asg_fork_id_t fork,
asg_record_id_t start_record,
asg_size_t recordcount,
- asg_size_t recordlen,
asg_flags_t flags,
asg_version_t version_condition,
- asg_version_t new_version,
+ asg_version_t * new_version,
asg_size_t * transferred)
{
return ASG_SUCCESS;
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. 7862829659b1f8844f2e058aad99d85f38c2bfcf
by noreply@mcs.anl.gov 25 Apr '14
by noreply@mcs.anl.gov 25 Apr '14
25 Apr '14
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, trac-277-bulkio has been updated
via 7862829659b1f8844f2e058aad99d85f38c2bfcf (commit)
from f5385cd8fd088bf4a146c823f31e38cb559dd05e (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 7862829659b1f8844f2e058aad99d85f38c2bfcf
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Apr 25 14:29:20 2014 -0400
added "big" triton-cp test
- copies 3-way replicated, 128 MB object on and off of 4 server
configuration using 32 MB I/O operations and 4 MB pipeline buffer size
- doesn't work reliably yet
-----------------------------------------------------------------------
Summary of changes:
code/tests/Makefile.subdir | 1 +
code/tests/triton-cp-big.sh | 64 +++++++++++++++++++++++++++++++++++++++++++
2 files changed, 65 insertions(+), 0 deletions(-)
create mode 100755 code/tests/triton-cp-big.sh
Diff of changes:
diff --git a/code/tests/Makefile.subdir b/code/tests/Makefile.subdir
index 8912e05..e95af09 100644
--- a/code/tests/Makefile.subdir
+++ b/code/tests/Makefile.subdir
@@ -18,6 +18,7 @@ TESTS += \
tests/triton-noop-multi-svr.sh \
tests/triton-touch.sh \
tests/triton-cp.sh \
+ tests/triton-cp-big.sh \
tests/triton-ls.sh \
tests/triton-rm.sh \
tests/triton-show-system-state.sh \
diff --git a/code/tests/triton-cp-big.sh b/code/tests/triton-cp-big.sh
new file mode 100755
index 0000000..60a76d7
--- /dev/null
+++ b/code/tests/triton-cp-big.sh
@@ -0,0 +1,64 @@
+#!/bin/bash
+
+if [ -z $srcdir ]; then
+ echo srcdir variable not set.
+ exit 1
+fi
+source $srcdir/tests/test-util.sh
+
+# start 4 servers with 15 second wait, 120s timeout
+test_start_servers 4 15 120
+
+# actual test case
+#####################
+
+# generate data file. We start with 1000000 bytes of random data and
+# then create a repeating pattern. The original random data is aligned on
+# power of 10 rather than power of 2 boundary on purpose to make sure that
+# we can detect errors that might occur on pipeline buffer boundaries.
+#
+dd if=/dev/urandom of=/tmp/128-$$.dat bs=1000000 count=1
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=1 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=2 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=4 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=8 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=16 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=32 conv=notrunc oflag=append
+dd if=/tmp/128-$$.dat of=/tmp/128-$$.dat bs=1M count=64 conv=notrunc oflag=append
+
+run_to 90 src/admin-tools/triton-cp $svr1 posix:/tmp/128-$$.dat ::128 3
+if [ $? -ne 0 ]; then
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+run_to 90 src/admin-tools/triton-cp $svr1 ::128 posix:/tmp/128-$$.dat.check
+if [ $? -ne 0 ]; then
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+run_to 90 src/admin-tools/triton-rm --server $svr1 ::128
+if [ $? -ne 0 ]; then
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+diff /tmp/128-$$.dat /tmp/128-$$.dat.check
+if [ $? -ne 0 ]; then
+ echo ERROR: data copied in and out of Triton does not match.
+ run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+ wait
+ exit 1
+fi
+
+#####################
+
+# tear down
+run_to 60 src/admin-tools/triton-shutdown-all-servers $svr1 &> /dev/null
+
+wait
+exit 0
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. f5385cd8fd088bf4a146c823f31e38cb559dd05e
by noreply@mcs.anl.gov 25 Apr '14
by noreply@mcs.anl.gov 25 Apr '14
25 Apr '14
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, trac-277-bulkio has been updated
via f5385cd8fd088bf4a146c823f31e38cb559dd05e (commit)
via 5d8e5aaa8910c3ad52daf9be5c61a27cbc955003 (commit)
from 7842dc0aa74f5d6dc9a8ec5fc25a3f32c8067ba6 (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 f5385cd8fd088bf4a146c823f31e38cb559dd05e
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Apr 25 13:06:57 2014 -0400
bug fix
- pipelined ROSD write now works for larger examples
- needs further testing with data validation and replication
commit 5d8e5aaa8910c3ad52daf9be5c61a27cbc955003
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Apr 25 12:35:39 2014 -0400
minor cleanup
-----------------------------------------------------------------------
Summary of changes:
code/src/replicated-osd/rosd-write.ae | 20 ++++++++++----------
code/src/replicated-osd/rosd.ae | 16 +++++++++++++++-
2 files changed, 25 insertions(+), 11 deletions(-)
Diff of changes:
diff --git a/code/src/replicated-osd/rosd-write.ae b/code/src/replicated-osd/rosd-write.ae
index dcba9b6..5d89e95 100644
--- a/code/src/replicated-osd/rosd-write.ae
+++ b/code/src/replicated-osd/rosd-write.ae
@@ -109,7 +109,6 @@ static __blocking triton_ret_t rosd_write_pull_one_buffer(
uint128_t oid,
uint64_t fork,
int64_t size,
- int64_t offset,
int64_t remote_offset,
hg_bulk_t remote_bulk_handle)
{
@@ -153,13 +152,13 @@ static __blocking triton_ret_t rosd_write_run_iteration(
triton_rpc_rosd_write_in_t* in,
aesop_sem_t *sem,
int* size_remaining,
+ int* this_size,
int64_t *remote_offset,
int64_t *local_offset,
struct buffer_mgmt_token **token,
char** buffer,
int* more_flag)
{
- int this_size = 0;
int64_t this_local_offset = 0;
int64_t this_remote_offset = 0;
triton_ret_t tret;
@@ -178,16 +177,16 @@ static __blocking triton_ret_t rosd_write_run_iteration(
size_remaining,
remote_offset,
local_offset,
- &this_size,
+ this_size,
&this_local_offset,
&this_remote_offset,
token,
buffer);
/* TODO: error handling */
assert(!triton_is_error(tret));
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished setup_one_buffer(), size_remaining: %d, this_size: %d.\n", *size_remaining, this_size);
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished setup_one_buffer(), this_size: %d, this_local_offset: %ld, this_remote_offset %ld.\n", *this_size, this_local_offset, this_remote_offset);
- if(this_size > 0)
+ if(*this_size > 0)
{
/* perform buffer transfer */
triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling pull_one_buffer().\n");
@@ -196,8 +195,7 @@ static __blocking triton_ret_t rosd_write_run_iteration(
*buffer,
in->oid,
in->oid_fork,
- this_size,
- this_local_offset,
+ *this_size,
this_remote_offset,
in->bulk_handle);
/* TODO: error handling */
@@ -206,7 +204,7 @@ static __blocking triton_ret_t rosd_write_run_iteration(
triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch starting do_work().\n");
tret = rosd_write_do_work(next_addr, in->oid,
- in->oid_fork, *buffer, this_size, this_local_offset,
+ in->oid_fork, *buffer, *this_size, this_local_offset,
in->flags, in->replication_factor, in->txn_number,
my_position, from_client_flag);
/* TODO: error handling */
@@ -245,6 +243,7 @@ static __blocking triton_ret_t rosd_write_run_pipeline(
pprivate struct buffer_mgmt_token *token = NULL;
pprivate char* buffer = NULL;
pprivate int more_flag = 1;
+ pprivate int this_size = 0;
triton_debug(triton_dbg_rosd,
"rosd_write_run_pipeline() starting pwait, size_remaining: %d.\n",
@@ -258,8 +257,9 @@ static __blocking triton_ret_t rosd_write_run_pipeline(
do
{
- tret = rosd_write_run_iteration(src_addr, next_addr, my_position,
- from_client_flag, in, &sem, &size_remaining, &remote_offset,
+ tret = rosd_write_run_iteration(src_addr, next_addr,
+ my_position, from_client_flag, in, &sem,
+ &size_remaining, &this_size, &remote_offset,
&local_offset, &token, &buffer, &more_flag);
/* TODO: error handling */
assert(!triton_is_error(tret));
diff --git a/code/src/replicated-osd/rosd.ae b/code/src/replicated-osd/rosd.ae
index 121ed6b..f70c89e 100644
--- a/code/src/replicated-osd/rosd.ae
+++ b/code/src/replicated-osd/rosd.ae
@@ -409,10 +409,14 @@ __blocking triton_ret_t rosd_pipeline_setup_one_buffer(
else
return(TRITON_ERR_UNKNOWN);
}
-
+
+ triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() acquired semaphore with size_remaining %d\n", *size_remaining);
+
if((*size_remaining) == 0)
{
*this_size = 0; /* done */
+ *this_local_offset = 0;
+ *this_remote_offset = 0;
ret = aesop_sem_up(sem);
assert(ret == AE_SUCCESS);
return(TRITON_SUCCESS);
@@ -430,6 +434,14 @@ __blocking triton_ret_t rosd_pipeline_setup_one_buffer(
return(tret);
}
}
+ else
+ {
+ /* If we have a buffer, then it better be non-zero size and have a
+ * token associated with it
+ */
+ assert(*this_size > 0);
+ assert(*token);
+ }
/* figure out what region we are accessing now */
if(*this_size > *size_remaining)
@@ -442,6 +454,8 @@ __blocking triton_ret_t rosd_pipeline_setup_one_buffer(
*local_offset += *this_size;
*remote_offset += *this_size;
+ triton_debug(triton_dbg_rosd, "rosd_pipeline_setup_one_buffer() releasing semaphore with size_remaining %d\n", *size_remaining);
+
ret = aesop_sem_up(sem);
assert(ret == AE_SUCCESS);
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. 7842dc0aa74f5d6dc9a8ec5fc25a3f32c8067ba6
by noreply@mcs.anl.gov 25 Apr '14
by noreply@mcs.anl.gov 25 Apr '14
25 Apr '14
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, trac-277-bulkio has been updated
via 7842dc0aa74f5d6dc9a8ec5fc25a3f32c8067ba6 (commit)
via 3cd538544e5c8e223ea35991ae4fd39e298a1e6a (commit)
from c8118f1185e5bc26a8252131839bb1559c7ef8d8 (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 7842dc0aa74f5d6dc9a8ec5fc25a3f32c8067ba6
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Apr 25 12:05:44 2014 -0400
reintegrate some functions
commit 3cd538544e5c8e223ea35991ae4fd39e298a1e6a
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Fri Apr 25 11:40:00 2014 -0400
switch back to do/while loop
-----------------------------------------------------------------------
Summary of changes:
code/src/replicated-osd/rosd-write.ae | 65 ++++++++++++---------------------
1 files changed, 24 insertions(+), 41 deletions(-)
Diff of changes:
diff --git a/code/src/replicated-osd/rosd-write.ae b/code/src/replicated-osd/rosd-write.ae
index 575c0b4..dcba9b6 100644
--- a/code/src/replicated-osd/rosd-write.ae
+++ b/code/src/replicated-osd/rosd-write.ae
@@ -219,40 +219,6 @@ static __blocking triton_ret_t rosd_write_run_iteration(
return(TRITON_SUCCESS);
}
-static __blocking triton_ret_t rosd_write_run_branch(
- na_addr_t src_addr,
- na_addr_t next_addr,
- int my_position,
- int from_client_flag,
- triton_rpc_rosd_write_in_t* in,
- aesop_sem_t *sem,
- int* size_remaining,
- int64_t *remote_offset,
- int64_t *local_offset)
-{
- struct buffer_mgmt_token *token = NULL;
- char* buffer = NULL;
- triton_ret_t tret;
- int more_flag;
-
- triton_debug(triton_dbg_rosd, "rosd_write_run_branch() starting.\n");
-
- for(more_flag=1; more_flag==1;)
- {
- tret = rosd_write_run_iteration(src_addr, next_addr, my_position,
- from_client_flag, in, sem, size_remaining, remote_offset,
- local_offset, &token, &buffer, &more_flag);
- }
-
- if(token != NULL)
- {
- buffer_mgmt_free(token);
- }
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() done with pbranch.\n");
-
- return(TRITON_SUCCESS);
-}
-
static __blocking triton_ret_t rosd_write_run_pipeline(
na_addr_t src_addr,
na_addr_t next_addr,
@@ -276,18 +242,35 @@ static __blocking triton_ret_t rosd_write_run_pipeline(
{
pprivate int i = 0;
pprivate triton_ret_t tret;
+ pprivate struct buffer_mgmt_token *token = NULL;
+ pprivate char* buffer = NULL;
+ pprivate int more_flag = 1;
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() in pwait, size_remaining: %d.\n", size_remaining);
+ triton_debug(triton_dbg_rosd,
+ "rosd_write_run_pipeline() starting pwait, size_remaining: %d.\n",
+ size_remaining);
for(i=0; i<ROSD_WRITE_XFER_PIPELINE_DEPTH; i++)
{
pbranch
{
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() starting pbranch %d.\n", i);
- tret = rosd_write_run_branch(src_addr, next_addr,
- my_position, from_client_flag, in, &sem, &size_remaining,
- &remote_offset, &local_offset);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
+ triton_debug(triton_dbg_rosd,
+ "rosd_write_run_pipeline() starting pbranch %d.\n", i);
+
+ do
+ {
+ tret = rosd_write_run_iteration(src_addr, next_addr, my_position,
+ from_client_flag, in, &sem, &size_remaining, &remote_offset,
+ &local_offset, &token, &buffer, &more_flag);
+ /* TODO: error handling */
+ assert(!triton_is_error(tret));
+ }while(more_flag);
+
+ if(token != NULL)
+ {
+ buffer_mgmt_free(token);
+ }
+ triton_debug(triton_dbg_rosd,
+ "rosd_write_run_pipeline() finishing pbranch %d.\n", i);
}
}
}
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. c8118f1185e5bc26a8252131839bb1559c7ef8d8
by noreply@mcs.anl.gov 24 Apr '14
by noreply@mcs.anl.gov 24 Apr '14
24 Apr '14
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, trac-277-bulkio has been updated
via c8118f1185e5bc26a8252131839bb1559c7ef8d8 (commit)
from d6d42eb12f4c3eee62758701a9a948afe1b82688 (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 c8118f1185e5bc26a8252131839bb1559c7ef8d8
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 21:20:56 2014 -0400
boost triton-cp buffer size to 32M
- needs debugging for write pipelining path: data written to server is
incomplete when moving large files even though triton-cp reports
success
-----------------------------------------------------------------------
Summary of changes:
code/src/admin-tools/triton-cp.ae | 2 +-
1 files changed, 1 insertions(+), 1 deletions(-)
Diff of changes:
diff --git a/code/src/admin-tools/triton-cp.ae b/code/src/admin-tools/triton-cp.ae
index a27498b..97c812a 100644
--- a/code/src/admin-tools/triton-cp.ae
+++ b/code/src/admin-tools/triton-cp.ae
@@ -17,7 +17,7 @@
#include "src/replicated-osd/rosd.hae"
/* TODO: make this configurable */
-#define BUFFER_SZ (4*1024*1024)
+#define BUFFER_SZ (32*1024*1024)
enum obj_ref_type
{
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. d6d42eb12f4c3eee62758701a9a948afe1b82688
by noreply@mcs.anl.gov 24 Apr '14
by noreply@mcs.anl.gov 24 Apr '14
24 Apr '14
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, trac-277-bulkio has been updated
via d6d42eb12f4c3eee62758701a9a948afe1b82688 (commit)
from ba7178daa0efdea03245ebca1707fc68e59e1e19 (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 d6d42eb12f4c3eee62758701a9a948afe1b82688
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 20:59:24 2014 -0400
reorganize rosd write further
- works for basic examples now
- need to follow up on an Aesop bug
- code needs to be cleaned up
-----------------------------------------------------------------------
Summary of changes:
code/src/replicated-osd/rosd-write.ae | 84 ++++++++++++++++++---------------
1 files changed, 46 insertions(+), 38 deletions(-)
Diff of changes:
diff --git a/code/src/replicated-osd/rosd-write.ae b/code/src/replicated-osd/rosd-write.ae
index e2dd06d..575c0b4 100644
--- a/code/src/replicated-osd/rosd-write.ae
+++ b/code/src/replicated-osd/rosd-write.ae
@@ -215,11 +215,6 @@ static __blocking triton_ret_t rosd_write_run_iteration(
*more_flag = 1;
}
- else
- {
- /* TODO: big hack! */
- aesop_timer(2);
- }
return(TRITON_SUCCESS);
}
@@ -258,6 +253,50 @@ static __blocking triton_ret_t rosd_write_run_branch(
return(TRITON_SUCCESS);
}
+static __blocking triton_ret_t rosd_write_run_pipeline(
+ na_addr_t src_addr,
+ na_addr_t next_addr,
+ int my_position,
+ int from_client_flag,
+ triton_rpc_rosd_write_in_t* in
+)
+{
+ /* TODO: 64 bit? */
+ int size_remaining = 0;
+ int64_t remote_offset = 0;
+ int64_t local_offset = 0;
+ aesop_sem_t sem;
+
+ size_remaining = in->size;
+ local_offset = in->offset;
+ remote_offset = 0;
+ aesop_sem_init(&sem, 1);
+
+ pwait
+ {
+ pprivate int i = 0;
+ pprivate triton_ret_t tret;
+
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() in pwait, size_remaining: %d.\n", size_remaining);
+ for(i=0; i<ROSD_WRITE_XFER_PIPELINE_DEPTH; i++)
+ {
+ pbranch
+ {
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() starting pbranch %d.\n", i);
+ tret = rosd_write_run_branch(src_addr, next_addr,
+ my_position, from_client_flag, in, &sem, &size_remaining,
+ &remote_offset, &local_offset);
+ /* TODO: error handling */
+ assert(!triton_is_error(tret));
+ }
+ }
+ }
+ aesop_sem_destroy(&sem);
+
+ return(TRITON_SUCCESS);
+}
+
+
static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
{
triton_rpc_rosd_write_out_t out;
@@ -273,11 +312,6 @@ static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
char* buffer_offsets[1];
struct txn_nr_cache_entry* entry_p = NULL;
na_addr_t src_addr;
- /* TODO: 64 bit? */
- int size_remaining = 0;
- int64_t remote_offset = 0;
- int64_t local_offset = 0;
- aesop_sem_t sem;
triton_debug(triton_dbg_rosd, "Called triton_rpc_rosd_write()\n");
@@ -381,40 +415,14 @@ static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
next_addr = addr_array[my_position+1];
}
- out.tret = TRITON_SUCCESS;
-
- size_remaining = in.size;
- local_offset = in.offset;
- remote_offset = 0;
- aesop_sem_init(&sem, 1);
-
- pwait
- {
- pprivate int i = 0;
- pprivate triton_ret_t tret;
-
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() in pwait, size_remaining: %d.\n", size_remaining);
- for(i=0; i<ROSD_WRITE_XFER_PIPELINE_DEPTH; i++)
- {
- pbranch
- {
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() starting pbranch %d.\n", i);
- tret = rosd_write_run_branch(src_addr, next_addr,
- my_position, from_client_flag, &in, &sem, &size_remaining,
- &remote_offset, &local_offset);
- /* TODO: error handling */
- assert(!triton_is_error(tret));
- }
- }
- }
- triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() done with pwait.\n");
+ out.tret = rosd_write_run_pipeline(src_addr, next_addr, my_position,
+ from_client_flag, &in);
if(entry_p)
txn_nr_cache_put(entry_p);
triton_mercury_start_output(handle, &out);
free(addr_array);
- aesop_sem_destroy(&sem);
return(TRITON_SUCCESS);
}
hooks/post-receive
--
1
0
branch, trac-277-bulkio, updated. ba7178daa0efdea03245ebca1707fc68e59e1e19
by noreply@mcs.anl.gov 24 Apr '14
by noreply@mcs.anl.gov 24 Apr '14
24 Apr '14
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, trac-277-bulkio has been updated
via ba7178daa0efdea03245ebca1707fc68e59e1e19 (commit)
via 5b111a97c8df0ccab0287c94337f89004602c9e7 (commit)
via ecfda6380e6119eb1c7739e7f54fd5795d24d175 (commit)
via 035ae077293f5c81c31eecd028d22ce5d7561fab (commit)
via 59d9c8b3159722ba5d9e1b3a4a04ebc9733c55c1 (commit)
via 2951c7427e96f3caedfb51e2ac89b118f7770d62 (commit)
via d1fd10f7a6caea9cd6ff47f98bdd286b2dd72ee2 (commit)
via 58dce24dbe05a8fe29c05db3be83386e96bf880f (commit)
from ed264ff8dfe6584f4c0c6d81bbb9a2cc8cf91456 (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 ba7178daa0efdea03245ebca1707fc68e59e1e19
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 17:25:03 2014 -0400
refactoring
- this isn't really the right organization, just trying to isolate a bug
- still crashes on write
commit 5b111a97c8df0ccab0287c94337f89004602c9e7
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 16:08:17 2014 -0400
more debugging msgs
commit ecfda6380e6119eb1c7739e7f54fd5795d24d175
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 16:05:45 2014 -0400
more debugging msgs
commit 035ae077293f5c81c31eecd028d22ce5d7561fab
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 16:01:38 2014 -0400
debugging msgs
commit 59d9c8b3159722ba5d9e1b3a4a04ebc9733c55c1
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 15:51:23 2014 -0400
draft of write pipelining
doesn't work yet
commit 2951c7427e96f3caedfb51e2ac89b118f7770d62
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 15:31:57 2014 -0400
subroutine for rdma pipeline xfer
commit d1fd10f7a6caea9cd6ff47f98bdd286b2dd72ee2
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 15:22:27 2014 -0400
move setup_buffer function to shared location
commit 58dce24dbe05a8fe29c05db3be83386e96bf880f
Author: Phil Carns <carns(a)mcs.anl.gov>
Date: Thu Apr 24 14:50:16 2014 -0400
reorder bulk buffer mgmt steps a little
-----------------------------------------------------------------------
Summary of changes:
code/src/replicated-osd/rosd-read.ae | 65 +---------
code/src/replicated-osd/rosd-write.ae | 239 ++++++++++++++++++++++++++-------
code/src/replicated-osd/rosd.ae | 63 +++++++++
code/src/replicated-osd/rosd.hae | 14 ++
4 files changed, 269 insertions(+), 112 deletions(-)
Diff of changes:
diff --git a/code/src/replicated-osd/rosd-read.ae b/code/src/replicated-osd/rosd-read.ae
index d4c4b1e..022b033 100644
--- a/code/src/replicated-osd/rosd-read.ae
+++ b/code/src/replicated-osd/rosd-read.ae
@@ -124,68 +124,6 @@ __blocking triton_ret_t remote_triton_rpc_rosd_read(
}
-static __blocking triton_ret_t rosd_read_setup_one_buffer(
- aesop_sem_t *sem,
- int* size_remaining,
- int64_t *remote_offset,
- int64_t *local_offset,
- int* this_size,
- int64_t* this_local_offset,
- int64_t* this_remote_offset,
- struct buffer_mgmt_token **token,
- char** buffer
- )
-{
- int ret;
- triton_ret_t tret;
-
- ret = aesop_sem_down(sem);
- if(ret != AE_SUCCESS)
- {
- if(ret == AE_ERR_CANCELLED)
- return(TRITON_ERR_CANCELED);
- else
- return(TRITON_ERR_UNKNOWN);
- }
-
- if((*size_remaining) == 0)
- {
- *this_size = 0; /* done */
- ret = aesop_sem_up(sem);
- assert(ret == AE_SUCCESS);
- return(TRITON_SUCCESS);
- }
-
- /* do we have a buffer for this pbranch yet? */
- if((*buffer) == NULL)
- {
- tret = buffer_mgmt_alloc(rosd_read_buffers,
- *size_remaining, this_size, buffer, token);
- if(triton_is_error(tret))
- {
- ret = aesop_sem_up(sem);
- assert(ret == AE_SUCCESS);
- return(tret);
- }
- }
-
- /* figure out what region we are accessing now */
- if(*this_size > *size_remaining)
- *this_size = *size_remaining;
- *this_local_offset = *local_offset;
- *this_remote_offset = *remote_offset;
-
- /* update offsets/sizes for overall xfer */
- *size_remaining -= *this_size;
- *local_offset += *this_size;
- *remote_offset += *this_size;
-
- ret = aesop_sem_up(sem);
- assert(ret == AE_SUCCESS);
-
- return(TRITON_SUCCESS);
-}
-
static __blocking triton_ret_t rosd_read_xfer_one_buffer(
na_addr_t src_addr,
char* tmp_buffer,
@@ -328,7 +266,8 @@ static __blocking triton_ret_t triton_rpc_rosd_read(hg_handle_t handle)
* using a semaphore
*/
//triton_debug(triton_dbg_rosd, "triton_rpc_rosd_read() pbranch %d calling setup_one_buffer().\n", i);
- tret = rosd_read_setup_one_buffer(
+ tret = rosd_pipeline_setup_one_buffer(
+ rosd_read_buffers,
&sem,
&size_remaining,
&remote_offset,
diff --git a/code/src/replicated-osd/rosd-write.ae b/code/src/replicated-osd/rosd-write.ae
index f0efe74..e2dd06d 100644
--- a/code/src/replicated-osd/rosd-write.ae
+++ b/code/src/replicated-osd/rosd-write.ae
@@ -17,6 +17,12 @@
#include "src/system-state/system-state.hae"
#include "src/transactional-osd/transactional-osd.hae"
+/* TODO: make this configurable */
+/* this is the number of concurrent buffers/transfers that the ROSD will
+ * keep in flight for a single write operation
+ */
+#define ROSD_WRITE_XFER_PIPELINE_DEPTH 4
+
/* Mercury RPC structures for rosd_write */
MERCURY_GEN_PROC(triton_rpc_rosd_write_out_t, ((triton_ret_t)(tret)))
MERCURY_GEN_PROC(triton_rpc_rosd_write_in_t,
@@ -31,6 +37,7 @@ MERCURY_GEN_PROC(triton_rpc_rosd_write_in_t,
((hg_bulk_t)(bulk_handle)))
+extern struct buffer_mgmt_instance* rosd_write_buffers;
static hg_id_t rpc_rosd_write_id;
static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle);
static int triton_rpc_rosd_write_handler(hg_handle_t handle);
@@ -96,7 +103,161 @@ __blocking triton_ret_t remote_triton_rpc_rosd_write(
}
-/* TODO: refactor this function; there is a lot going on in here now */
+static __blocking triton_ret_t rosd_write_pull_one_buffer(
+ na_addr_t src_addr,
+ char* tmp_buffer,
+ uint128_t oid,
+ uint64_t fork,
+ int64_t size,
+ int64_t offset,
+ int64_t remote_offset,
+ hg_bulk_t remote_bulk_handle)
+{
+ triton_ret_t tret;
+ hg_bulk_t bulk_handle = HG_BULK_NULL;
+ hg_bulk_request_t bulk_request;
+ int ret;
+
+ ret = HG_Bulk_handle_create(tmp_buffer, size, HG_BULK_READWRITE, &bulk_handle);
+ if(ret != HG_SUCCESS)
+ {
+ triton_error_msg("HG_Bulk_handle_create() failure.\n");
+ return(TRITON_ERR_NOMEM);
+ }
+
+ ret = HG_Bulk_read(src_addr, remote_bulk_handle, remote_offset,
+ bulk_handle, 0, size, &bulk_request);
+ if(ret != HG_SUCCESS)
+ {
+ triton_error_msg("HG_Bulk_read() failure.\n");
+ HG_Bulk_handle_free(bulk_handle);
+ return(TRITON_ERR_UNKNOWN);
+ }
+
+ tret = triton_mercury_bulk_wait(bulk_request);
+ if(triton_is_error(tret))
+ {
+ triton_error_msg("triton_mercury_bulk_wait() failure.\n");
+ }
+
+ HG_Bulk_handle_free(bulk_handle);
+
+ return(tret);
+}
+
+static __blocking triton_ret_t rosd_write_run_iteration(
+ na_addr_t src_addr,
+ na_addr_t next_addr,
+ int my_position,
+ int from_client_flag,
+ triton_rpc_rosd_write_in_t* in,
+ aesop_sem_t *sem,
+ int* size_remaining,
+ int64_t *remote_offset,
+ int64_t *local_offset,
+ struct buffer_mgmt_token **token,
+ char** buffer,
+ int* more_flag)
+{
+ int this_size = 0;
+ int64_t this_local_offset = 0;
+ int64_t this_remote_offset = 0;
+ triton_ret_t tret;
+
+ *more_flag = 0;
+
+ /* calculate how much to transfer in this step */
+
+ /* note that this function protects shared variables
+ * using a semaphore
+ */
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling setup_one_buffer().\n");
+ tret = rosd_pipeline_setup_one_buffer(
+ rosd_write_buffers,
+ sem,
+ size_remaining,
+ remote_offset,
+ local_offset,
+ &this_size,
+ &this_local_offset,
+ &this_remote_offset,
+ token,
+ buffer);
+ /* TODO: error handling */
+ assert(!triton_is_error(tret));
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished setup_one_buffer(), size_remaining: %d, this_size: %d.\n", *size_remaining, this_size);
+
+ if(this_size > 0)
+ {
+ /* perform buffer transfer */
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch calling pull_one_buffer().\n");
+ tret = rosd_write_pull_one_buffer(
+ src_addr,
+ *buffer,
+ in->oid,
+ in->oid_fork,
+ this_size,
+ this_local_offset,
+ this_remote_offset,
+ in->bulk_handle);
+ /* TODO: error handling */
+ assert(!triton_is_error(tret));
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished pull_one_buffer().\n");
+
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch starting do_work().\n");
+ tret = rosd_write_do_work(next_addr, in->oid,
+ in->oid_fork, *buffer, this_size, this_local_offset,
+ in->flags, in->replication_factor, in->txn_number,
+ my_position, from_client_flag);
+ /* TODO: error handling */
+ assert(!triton_is_error(tret));
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() pbranch finished do_work().\n");
+
+ *more_flag = 1;
+ }
+ else
+ {
+ /* TODO: big hack! */
+ aesop_timer(2);
+ }
+
+ return(TRITON_SUCCESS);
+}
+
+static __blocking triton_ret_t rosd_write_run_branch(
+ na_addr_t src_addr,
+ na_addr_t next_addr,
+ int my_position,
+ int from_client_flag,
+ triton_rpc_rosd_write_in_t* in,
+ aesop_sem_t *sem,
+ int* size_remaining,
+ int64_t *remote_offset,
+ int64_t *local_offset)
+{
+ struct buffer_mgmt_token *token = NULL;
+ char* buffer = NULL;
+ triton_ret_t tret;
+ int more_flag;
+
+ triton_debug(triton_dbg_rosd, "rosd_write_run_branch() starting.\n");
+
+ for(more_flag=1; more_flag==1;)
+ {
+ tret = rosd_write_run_iteration(src_addr, next_addr, my_position,
+ from_client_flag, in, sem, size_remaining, remote_offset,
+ local_offset, &token, &buffer, &more_flag);
+ }
+
+ if(token != NULL)
+ {
+ buffer_mgmt_free(token);
+ }
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() done with pbranch.\n");
+
+ return(TRITON_SUCCESS);
+}
+
static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
{
triton_rpc_rosd_write_out_t out;
@@ -105,16 +266,18 @@ static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
na_addr_t next_addr;
na_addr_t* addr_array;
int my_position;
+ int from_client_flag = 0;
int64_t obj_offset;
int64_t size;
int64_t out_size;
char* buffer_offsets[1];
- int from_client_flag = 0;
struct txn_nr_cache_entry* entry_p = NULL;
- char* tmp_buffer;
- hg_bulk_t bulk_handle = HG_BULK_NULL;
na_addr_t src_addr;
- hg_bulk_request_t bulk_request;
+ /* TODO: 64 bit? */
+ int size_remaining = 0;
+ int64_t remote_offset = 0;
+ int64_t local_offset = 0;
+ aesop_sem_t sem;
triton_debug(triton_dbg_rosd, "Called triton_rpc_rosd_write()\n");
@@ -218,62 +381,40 @@ static __blocking triton_ret_t triton_rpc_rosd_write(hg_handle_t handle)
next_addr = addr_array[my_position+1];
}
- /* TODO: buffer management */
- /* TODO: pipelining */
- if(in.size > 0)
- {
- tmp_buffer = malloc(in.size);
- if(!tmp_buffer)
- {
- out.tret = TRITON_ERR_NOMEM;
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
- }
+ out.tret = TRITON_SUCCESS;
- ret = HG_Bulk_handle_create(tmp_buffer, in.size, HG_BULK_READWRITE, &bulk_handle);
- if(ret != HG_SUCCESS)
- {
- triton_error_msg("HG_Bulk_handle_create() failure.\n");
- out.tret = TRITON_ERR_NOMEM;
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
- }
+ size_remaining = in.size;
+ local_offset = in.offset;
+ remote_offset = 0;
+ aesop_sem_init(&sem, 1);
- ret = HG_Bulk_read(src_addr, in.bulk_handle, 0, bulk_handle,
- 0, in.size, &bulk_request);
- if(ret != HG_SUCCESS)
- {
- triton_error_msg("HG_Bulk_write() failure.\n");
- out.tret = TRITON_ERR_UNKNOWN;
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
- }
+ pwait
+ {
+ pprivate int i = 0;
+ pprivate triton_ret_t tret;
- out.tret = triton_mercury_bulk_wait(bulk_request);
- if(triton_is_error(out.tret))
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() in pwait, size_remaining: %d.\n", size_remaining);
+ for(i=0; i<ROSD_WRITE_XFER_PIPELINE_DEPTH; i++)
{
- triton_error_msg("triton_mercury_bulk_wait() failure.\n");
- free(addr_array);
- triton_mercury_start_output(handle, &out);
- return(TRITON_SUCCESS);
+ pbranch
+ {
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() starting pbranch %d.\n", i);
+ tret = rosd_write_run_branch(src_addr, next_addr,
+ my_position, from_client_flag, &in, &sem, &size_remaining,
+ &remote_offset, &local_offset);
+ /* TODO: error handling */
+ assert(!triton_is_error(tret));
+ }
}
}
-
- out.tret = rosd_write_do_work(next_addr, in.oid, in.oid_fork,
- tmp_buffer, in.size, in.offset, in.flags,
- in.replication_factor, in.txn_number, my_position, from_client_flag);
+ triton_debug(triton_dbg_rosd, "triton_rpc_rosd_write() done with pwait.\n");
if(entry_p)
txn_nr_cache_put(entry_p);
triton_mercury_start_output(handle, &out);
- if(in.size > 0)
- HG_Bulk_handle_free(bulk_handle);
- free(tmp_buffer);
free(addr_array);
+ aesop_sem_destroy(&sem);
return(TRITON_SUCCESS);
}
diff --git a/code/src/replicated-osd/rosd.ae b/code/src/replicated-osd/rosd.ae
index 4017aba..121ed6b 100644
--- a/code/src/replicated-osd/rosd.ae
+++ b/code/src/replicated-osd/rosd.ae
@@ -385,6 +385,69 @@ __blocking void trigger_server_fault(triton_ret_t tret)
return;
}
+__blocking triton_ret_t rosd_pipeline_setup_one_buffer(
+ struct buffer_mgmt_instance* buffer_instance,
+ aesop_sem_t *sem,
+ int* size_remaining,
+ int64_t *remote_offset,
+ int64_t *local_offset,
+ int* this_size,
+ int64_t* this_local_offset,
+ int64_t* this_remote_offset,
+ struct buffer_mgmt_token **token,
+ char** buffer
+ )
+{
+ int ret;
+ triton_ret_t tret;
+
+ ret = aesop_sem_down(sem);
+ if(ret != AE_SUCCESS)
+ {
+ if(ret == AE_ERR_CANCELLED)
+ return(TRITON_ERR_CANCELED);
+ else
+ return(TRITON_ERR_UNKNOWN);
+ }
+
+ if((*size_remaining) == 0)
+ {
+ *this_size = 0; /* done */
+ ret = aesop_sem_up(sem);
+ assert(ret == AE_SUCCESS);
+ return(TRITON_SUCCESS);
+ }
+
+ /* do we have a buffer for this pbranch yet? */
+ if((*buffer) == NULL)
+ {
+ tret = buffer_mgmt_alloc(buffer_instance,
+ *size_remaining, this_size, buffer, token);
+ if(triton_is_error(tret))
+ {
+ ret = aesop_sem_up(sem);
+ assert(ret == AE_SUCCESS);
+ return(tret);
+ }
+ }
+
+ /* figure out what region we are accessing now */
+ if(*this_size > *size_remaining)
+ *this_size = *size_remaining;
+ *this_local_offset = *local_offset;
+ *this_remote_offset = *remote_offset;
+
+ /* update offsets/sizes for overall xfer */
+ *size_remaining -= *this_size;
+ *local_offset += *this_size;
+ *remote_offset += *this_size;
+
+ ret = aesop_sem_up(sem);
+ assert(ret == AE_SUCCESS);
+
+ return(TRITON_SUCCESS);
+}
+
/*
* Local Variables:
* c-basic-offset: 4
diff --git a/code/src/replicated-osd/rosd.hae b/code/src/replicated-osd/rosd.hae
index 44d3c17..da89f27 100644
--- a/code/src/replicated-osd/rosd.hae
+++ b/code/src/replicated-osd/rosd.hae
@@ -2,10 +2,12 @@
#define __ROSD_HAE__
#include <aesop/aesop.h>
+#include <aesop/sem.hae>
#include <mercury.h>
#include <triton-uint128.h>
#include "src/common/triton-error.h"
+#include "src/replicated-osd/buffer-mgmt.hae"
/* for requests that should be fanned out from the master for replication */
#define ROSD_FLAG_FANOUT 1
@@ -62,6 +64,18 @@ __blocking triton_ret_t remote_triton_rpc_rosd_read(
int64_t* out_size,
uint32_t flags);
+__blocking triton_ret_t rosd_pipeline_setup_one_buffer(
+ struct buffer_mgmt_instance* buffer_instance,
+ aesop_sem_t *sem,
+ int* size_remaining,
+ int64_t *remote_offset,
+ int64_t *local_offset,
+ int* this_size,
+ int64_t* this_local_offset,
+ int64_t* this_remote_offset,
+ struct buffer_mgmt_token **token,
+ char** buffer
+ );
#endif /* __ROSD_HAE */
hooks/post-receive
--
1
0