Grayskull Repository branch, master, updated. git-migration-365-g0ab2aaf
A ref change was pushed to the repository containing the project "Grayskull Repository". The branch, master has been updated via 0ab2aafa5ccb295dad1a766c376072c0383620e2 (commit) via d98a47022f1523d537d7803f1360b226ce931d3f (commit) via 720bd0a4c7664584bc96bae602e1aa4de9bd6b72 (commit) via 4d7c42c8580f692eb51074d8b69588d5ace20680 (commit) via 01d213ca4b01136efcf02bded9d515da63e427b7 (commit) from b7520255a7b588bcbb4eabfafd9499e674338cdb (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 0ab2aafa5ccb295dad1a766c376072c0383620e2 Author: Phil Carns <[email protected]> Date: Fri Jan 29 13:21:14 2010 -0500 test auto commit and fix memory leak commit d98a47022f1523d537d7803f1360b226ce931d3f Author: Phil Carns <[email protected]> Date: Fri Jan 29 13:20:46 2010 -0500 fix to earlier memset in opcache commit 720bd0a4c7664584bc96bae602e1aa4de9bd6b72 Author: Phil Carns <[email protected]> Date: Fri Jan 29 12:46:02 2010 -0500 flag to enable auto-commit of individual writes Added a little infrastructure to combine multiple COSD operations into compound operations, used to implement COSD_FLAG_AUTO_TXN. Untested. commit 4d7c42c8580f692eb51074d8b69588d5ace20680 Author: Phil Carns <[email protected]> Date: Fri Jan 29 12:45:08 2010 -0500 zero out op structures for convenience commit 01d213ca4b01136efcf02bded9d515da63e427b7 Author: Phil Carns <[email protected]> Date: Fri Jan 29 10:29:40 2010 -0500 move op_worker function to be resource-specific ----------------------------------------------------------------------- Summary of changes: code/src/gsl/common/gs-opcache.c | 4 + code/src/gsl/include/gs-op.h | 1 - .../gsl/resources/cosd-prototype/cosd-prototype.c | 573 ++++++++++++-------- .../resources/cosd-prototype/cosd-prototype.gsh | 1 + .../src/gsl/resources/cosd-prototype/test/cosd1.gs | 21 + code/src/gsl/resources/storage/gs-storage.c | 21 +- 6 files changed, 386 insertions(+), 235 deletions(-) Diff of changes: diff --git a/code/src/gsl/common/gs-opcache.c b/code/src/gsl/common/gs-opcache.c index 63ab731..724e7a8 100644 --- a/code/src/gsl/common/gs-opcache.c +++ b/code/src/gsl/common/gs-opcache.c @@ -108,6 +108,8 @@ struct gs_op *gs_opcache_get(gs_opcache_t cache) op = (struct gs_op *)(((char *)cache->array[aind]) + (count * cache->typesize) + cache->member_offset); + memset((void*)((unsigned long)op - cache->member_offset), 0, + cache->typesize); op->cache_id = (aind << 25) | count; gs_op_clear(op); ++cache->count; @@ -115,6 +117,8 @@ struct gs_op *gs_opcache_get(gs_opcache_t cache) else { op = gs_oplist_pop(&cache->free_list); + memset((void*)((unsigned long)op - cache->member_offset), 0, + cache->typesize); } gs_mutex_unlock(&cache->mutex); return op; diff --git a/code/src/gsl/include/gs-op.h b/code/src/gsl/include/gs-op.h index 06515e3..a4debd7 100644 --- a/code/src/gsl/include/gs-op.h +++ b/code/src/gsl/include/gs-op.h @@ -7,7 +7,6 @@ struct gs_op { void (*callback)(void *ptr, int ret); - int (*op_worker)(struct gs_op* op); void *user_ptr; gs_hints_t hints; gs_context_t ctx; diff --git a/code/src/gsl/resources/cosd-prototype/cosd-prototype.c b/code/src/gsl/resources/cosd-prototype/cosd-prototype.c index 54b8c8b..2837ebd 100644 --- a/code/src/gsl/resources/cosd-prototype/cosd-prototype.c +++ b/code/src/gsl/resources/cosd-prototype/cosd-prototype.c @@ -96,23 +96,6 @@ static gs_mutex_t global_log_mutex = GS_MUTEX_INITIALIZER; static int gs_cosd_resource_id; static enum progress_mode gs_cosd_progress_mode = GS_PROG_NONE; -static int merge_logical_map(struct gs_list_link* list1_in, int - list1_count, struct gs_list_link* list2_in, int list2_count, - struct gs_list_link* list_out, void** free_ptr); -static int compare_log_map_key(DB * dbp, const DBT * a, const DBT * b); -static int compare_missing_version(DB * dbp, const DBT * a, const DBT * b); -static int compare_uint64(DB * dbp, const DBT * a, const DBT * b); -static int cmp_lme(const void *p1, const void *p2); -static int advance_listio_ptrs( - char** mem_offsets, - int64_t* mem_sizes, - int mem_count, - int* mem_index, - int64_t* obj_offsets, - int64_t* obj_sizes, - int obj_count, - int* obj_index, - int64_t size); /* counter used by poll function in thread-per-op case to know if ops have * finished since last poll @@ -165,9 +148,12 @@ struct missing_version uint64_t version; }; -struct cosd_op +struct cosd_op; + +struct cosd_work { - union { + union + { struct create_op{ uint64_t requested_oid; uint64_t* out_oid; @@ -204,30 +190,93 @@ struct cosd_op uint64_t oid; uint64_t *version; } get_version; - } u; + }u; + + int (*op_worker)(struct gs_op* op); + void (*cleanup_fn)(struct cosd_work* work); +}; + +struct cosd_op +{ + struct cosd_work work; + struct cosd_work* work_array; + int work_array_count; gs_op_id_t op_id; struct gs_op op; int error_code; - void (*cleanup_fn)(struct cosd_op* c_op); }; +static int merge_logical_map(struct gs_list_link* list1_in, int + list1_count, struct gs_list_link* list2_in, int list2_count, + struct gs_list_link* list_out, void** free_ptr); +static int compare_log_map_key(DB * dbp, const DBT * a, const DBT * b); +static int compare_missing_version(DB * dbp, const DBT * a, const DBT * b); +static int compare_uint64(DB * dbp, const DBT * a, const DBT * b); +static int cmp_lme(const void *p1, const void *p2); +static int advance_listio_ptrs( + char** mem_offsets, + int64_t* mem_sizes, + int mem_count, + int* mem_index, + int64_t* obj_offsets, + int64_t* obj_sizes, + int obj_count, + int* obj_index, + int64_t size); +static int cosd_write_set_work( + struct cosd_work* work, + uint64_t oid, + uint64_t fork, + uint64_t txn_number, + char** mem_offsets, + int64_t* mem_sizes, + int mem_count, + int64_t* obj_offsets, + int64_t* obj_sizes, + int obj_count, + int flags); +static int cosd_txn_close_set_work( + struct cosd_work* work, + uint64_t oid, + uint64_t fork, + uint64_t txn_number); + static void* thread_fn(void* foo) { struct gs_op* op = foo; struct cosd_op* c_op = gs_op_entry(op, struct cosd_op, op); + int i = 0; gs_mutex_lock(&cosd_mutex); /* pull off of list */ gs_oplist_del(op, &cosd_oplist); gs_mutex_unlock(&cosd_mutex); - /* spin on op worker */ - while(op->op_worker(op) != 1); + if(c_op->work_array_count == 0) + { + /* normal operation */ + /* spin on op worker */ + while(c_op->work.op_worker(op) != 1); - /* call cleanup function if present */ - if(c_op->cleanup_fn) + /* call cleanup function if present */ + if(c_op->work.cleanup_fn) + c_op->work.cleanup_fn(&c_op->work); + } + else { - c_op->cleanup_fn(c_op); + /* compound operation */ + while(i<c_op->work_array_count && c_op->error_code == 0) + { + c_op->work = c_op->work_array[i]; + while(c_op->work.op_worker(op) != 1); + i++; + } + for(i=0; i<c_op->work_array_count; i++) + { + if(c_op->work_array[i].cleanup_fn) + c_op->work_array[i].cleanup_fn(&c_op->work_array[i]); + } + free(c_op->work_array); } /* trigger completion of the operation */ @@ -353,7 +402,7 @@ static int gs_cosd_poll(gs_context_t context, int millisecs) struct gs_op *gop; struct cosd_op *c_op; struct timespec ts_sleep; - int ret; + int i=0; /* NOTE: just servicing one op per call right now */ @@ -383,25 +432,35 @@ static int gs_cosd_poll(gs_context_t context, int millisecs) c_op = gs_op_entry(gop, struct cosd_op, op); - ret = gop->op_worker(gop); - if(ret == 1) + if(c_op->work_array_count == 0) { - /* call cleanup function if present */ - if(c_op->cleanup_fn) - { - c_op->cleanup_fn(c_op); - } + /* normal operation */ + /* spin on op worker */ + while(c_op->work.op_worker(gop) != 1); - /* done */ - gs_opcache_complete_op(cosd_opcache, gop, c_op->error_code); + /* call cleanup function if present */ + if(c_op->work.cleanup_fn) + c_op->work.cleanup_fn(&c_op->work); } else { - assert(ret == 0); /* only 0 and 1 allowed? */ - /* not done */ - gs_oplist_add(gop, &cosd_oplist); + /* compound operation */ + while(i<c_op->work_array_count && c_op->error_code == 0) + { + c_op->work = c_op->work_array[i]; + while(c_op->work.op_worker(gop) != 1); + i++; + } + for(i=0; i<c_op->work_array_count; i++) + { + if(c_op->work_array[i].cleanup_fn) + c_op->work_array[i].cleanup_fn(&c_op->work_array[i]); + } + free(c_op->work_array); } + gs_opcache_complete_op(cosd_opcache, gop, c_op->error_code); + /* gs_mutex_unlock(&cosd_mutex); */ return 0; @@ -668,16 +727,16 @@ int gs_cosd_finalize(void) * * cleans up memory after a read operation completes */ -static void read_op_cleanup(struct cosd_op *c_op) +static void read_op_cleanup(struct cosd_work *work) { - if(c_op->u.read.mem_offsets) - free(c_op->u.read.mem_offsets); - if(c_op->u.read.mem_sizes) - free(c_op->u.read.mem_sizes); - if(c_op->u.read.obj_sizes) - free(c_op->u.read.obj_sizes); - if(c_op->u.read.obj_offsets) - free(c_op->u.read.obj_offsets); + if(work->u.read.mem_offsets) + free(work->u.read.mem_offsets); + if(work->u.read.mem_sizes) + free(work->u.read.mem_sizes); + if(work->u.read.obj_sizes) + free(work->u.read.obj_sizes); + if(work->u.read.obj_offsets) + free(work->u.read.obj_offsets); return; } @@ -686,16 +745,16 @@ static void read_op_cleanup(struct cosd_op *c_op) * * cleans up memory after a write operation completes */ -static void write_op_cleanup(struct cosd_op *c_op) +static void write_op_cleanup(struct cosd_work* work) { - if(c_op->u.write.mem_offsets) - free(c_op->u.write.mem_offsets); - if(c_op->u.write.mem_sizes) - free(c_op->u.write.mem_sizes); - if(c_op->u.write.obj_sizes) - free(c_op->u.write.obj_sizes); - if(c_op->u.write.obj_offsets) - free(c_op->u.write.obj_offsets); + if(work->u.write.mem_offsets) + free(work->u.write.mem_offsets); + if(work->u.write.mem_sizes) + free(work->u.write.mem_sizes); + if(work->u.write.obj_sizes) + free(work->u.write.obj_sizes); + if(work->u.write.obj_offsets) + free(work->u.write.obj_offsets); return; } @@ -725,18 +784,18 @@ static int read_op_worker(struct gs_op* op) assert(c_op); /* only support one object for now */ - assert(c_op->u.read.oid == 1); - lmk.oid = c_op->u.read.oid; - lmk.fork = c_op->u.read.fork; + assert(c_op->work.u.read.oid == 1); + lmk.oid = c_op->work.u.read.oid; + lmk.fork = c_op->work.u.read.fork; - assert(c_op->u.read.mem_count > 0); - assert(c_op->u.read.obj_count > 0); + assert(c_op->work.u.read.mem_count > 0); + assert(c_op->work.u.read.obj_count > 0); if(global_fd < 0) { /* open log file */ sprintf(log_name, "%s/%llu.dat", cosd_log_path, - llu(c_op->u.read.oid)); + llu(c_op->work.u.read.oid)); ret = open(log_name, O_RDWR|O_DIRECT|O_EXCL|O_NOATIME, S_IRUSR|S_IWUSR); if(ret < 0) @@ -788,7 +847,7 @@ read_op_retry: COSD_INIT_DBT(key, lmk); COSD_INIT_DBT(value, lme); - lmk.logical_offset_end = c_op->u.read.obj_offsets[obj_index] + 1; + lmk.logical_offset_end = c_op->work.u.read.obj_offsets[obj_index] + 1; ret = dbc_p->c_get(dbc_p, &key, &value, DB_SET_RANGE); if(ret == DB_NOTFOUND) @@ -842,27 +901,27 @@ read_op_retry: assert(0); } - if(lmk.oid != c_op->u.read.oid || lmk.fork != c_op->u.read.fork || - lme.logical_offset > c_op->u.read.obj_offsets[obj_index]) + if(lmk.oid != c_op->work.u.read.oid || lmk.fork != c_op->work.u.read.fork || + lme.logical_offset > c_op->work.u.read.obj_offsets[obj_index]) { int64_t amt_to_zero; int64_t amt_left; /* figure out how big this hole is */ - if(lmk.oid == c_op->u.read.oid && - lmk.fork == c_op->u.read.fork && - c_op->u.read.obj_sizes[obj_index] > (lme.logical_offset - - c_op->u.read.obj_offsets[obj_index])) + if(lmk.oid == c_op->work.u.read.oid && + lmk.fork == c_op->work.u.read.fork && + c_op->work.u.read.obj_sizes[obj_index] > (lme.logical_offset - + c_op->work.u.read.obj_offsets[obj_index])) { /* zero up to the next region in the log */ amt_to_zero = lme.logical_offset - - c_op->u.read.obj_offsets[obj_index]; + c_op->work.u.read.obj_offsets[obj_index]; } - else if(lmk.oid == c_op->u.read.oid && - lmk.fork == c_op->u.read.fork) + else if(lmk.oid == c_op->work.u.read.oid && + lmk.fork == c_op->work.u.read.fork) { /* zero this whole buffer */ - amt_to_zero = c_op->u.read.obj_sizes[obj_index]; + amt_to_zero = c_op->work.u.read.obj_sizes[obj_index]; } else { @@ -877,26 +936,26 @@ read_op_retry: { int64_t this_region; - if(amt_left > c_op->u.read.mem_sizes[mem_index]) - this_region = c_op->u.read.mem_sizes[mem_index]; + if(amt_left > c_op->work.u.read.mem_sizes[mem_index]) + this_region = c_op->work.u.read.mem_sizes[mem_index]; else this_region = amt_left; assert(this_region > 0); - memset(c_op->u.read.mem_offsets[mem_index], 0, this_region); + memset(c_op->work.u.read.mem_offsets[mem_index], 0, this_region); amt_left -= this_region; done = advance_listio_ptrs( - c_op->u.read.mem_offsets, - c_op->u.read.mem_sizes, - c_op->u.read.mem_count, + c_op->work.u.read.mem_offsets, + c_op->work.u.read.mem_sizes, + c_op->work.u.read.mem_count, &mem_index, - c_op->u.read.obj_offsets, - c_op->u.read.obj_sizes, - c_op->u.read.obj_count, + c_op->work.u.read.obj_offsets, + c_op->work.u.read.obj_sizes, + c_op->work.u.read.obj_count, &obj_index, this_region); - *c_op->u.read.out_size += this_region; + *c_op->work.u.read.out_size += this_region; /* keep going; there might be more data after the hole */ continue; @@ -906,17 +965,17 @@ read_op_retry: /* we have some data to read */ /* ignore portions of logical map beyond what we want */ - if(lme.logical_offset < c_op->u.read.obj_offsets[obj_index]) + if(lme.logical_offset < c_op->work.u.read.obj_offsets[obj_index]) { - int64_t diff = c_op->u.read.obj_offsets[obj_index] - + int64_t diff = c_op->work.u.read.obj_offsets[obj_index] - lme.logical_offset; lme.size -= diff; lme.logical_offset += diff; lme.log_offset += diff; } - if(lme.size > c_op->u.read.obj_sizes[obj_index]) + if(lme.size > c_op->work.u.read.obj_sizes[obj_index]) { - int64_t diff = lme.size - c_op->u.read.obj_sizes[obj_index]; + int64_t diff = lme.size - c_op->work.u.read.obj_sizes[obj_index]; lme.size -= diff; lme.logical_offset_end -= diff; } @@ -930,27 +989,27 @@ read_op_retry: */ if(lme.log_offset % DIRECT_ALIGN == 0 && lme.size % DIRECT_ALIGN == 0 && - ((unsigned long)c_op->u.read.mem_offsets[mem_index]) % DIRECT_ALIGN == 0 && - c_op->u.read.mem_sizes[mem_index] >= lme.size) + ((unsigned long)c_op->work.u.read.mem_offsets[mem_index]) % DIRECT_ALIGN == 0 && + c_op->work.u.read.mem_sizes[mem_index] >= lme.size) { int64_t pread_ret; - pread_ret = pread(global_fd, c_op->u.read.mem_offsets[mem_index], + pread_ret = pread(global_fd, c_op->work.u.read.mem_offsets[mem_index], lme.size, lme.log_offset); /* TODO: error handling */ assert(pread_ret == lme.size); /* advance pointers */ done = advance_listio_ptrs( - c_op->u.read.mem_offsets, - c_op->u.read.mem_sizes, - c_op->u.read.mem_count, + c_op->work.u.read.mem_offsets, + c_op->work.u.read.mem_sizes, + c_op->work.u.read.mem_count, &mem_index, - c_op->u.read.obj_offsets, - c_op->u.read.obj_sizes, - c_op->u.read.obj_count, + c_op->work.u.read.obj_offsets, + c_op->work.u.read.obj_sizes, + c_op->work.u.read.obj_count, &obj_index, lme.size); - *c_op->u.read.out_size += lme.size; + *c_op->work.u.read.out_size += lme.size; } else { @@ -998,27 +1057,27 @@ read_op_retry: copy_remaining = lme.size; while(copy_remaining) { - if(copy_remaining > c_op->u.read.mem_sizes[mem_index]) - copy_amt = c_op->u.read.mem_sizes[mem_index]; + if(copy_remaining > c_op->work.u.read.mem_sizes[mem_index]) + copy_amt = c_op->work.u.read.mem_sizes[mem_index]; else copy_amt = copy_remaining; - memcpy(c_op->u.read.mem_offsets[mem_index], copy_ptr, copy_amt); + memcpy(c_op->work.u.read.mem_offsets[mem_index], copy_ptr, copy_amt); copy_ptr += copy_amt; copy_remaining -= copy_amt; /* advance pointers */ done = advance_listio_ptrs( - c_op->u.read.mem_offsets, - c_op->u.read.mem_sizes, - c_op->u.read.mem_count, + c_op->work.u.read.mem_offsets, + c_op->work.u.read.mem_sizes, + c_op->work.u.read.mem_count, &mem_index, - c_op->u.read.obj_offsets, - c_op->u.read.obj_sizes, - c_op->u.read.obj_count, + c_op->work.u.read.obj_offsets, + c_op->work.u.read.obj_sizes, + c_op->work.u.read.obj_count, &obj_index, copy_amt); - *c_op->u.read.out_size += copy_amt; + *c_op->work.u.read.out_size += copy_amt; } } } @@ -1054,16 +1113,16 @@ static int write_op_worker(struct gs_op* op) assert(c_op); /* only support one object for now */ - assert(c_op->u.write.oid == 1); + assert(c_op->work.u.write.oid == 1); - assert(c_op->u.write.mem_count > 0); - assert(c_op->u.write.obj_count > 0); + assert(c_op->work.u.write.mem_count > 0); + assert(c_op->work.u.write.obj_count > 0); if(global_fd < 0) { /* open log file */ sprintf(log_name, "%s/%llu.dat", cosd_log_path, - llu(c_op->u.write.oid)); + llu(c_op->work.u.write.oid)); ret = open(log_name, O_RDWR|O_DIRECT|O_EXCL|O_NOATIME, S_IRUSR|S_IWUSR); if(ret < 0) @@ -1089,27 +1148,27 @@ static int write_op_worker(struct gs_op* op) } /* find the next chunk that is contiguous in memory and disk */ - mem_ptr = c_op->u.write.mem_offsets[mem_index]; - tmp_update->logical_offset = c_op->u.write.obj_offsets[obj_index]; - tmp_update->size = c_op->u.write.mem_sizes[mem_index]; + mem_ptr = c_op->work.u.write.mem_offsets[mem_index]; + tmp_update->logical_offset = c_op->work.u.write.obj_offsets[obj_index]; + tmp_update->size = c_op->work.u.write.mem_sizes[mem_index]; tmp_update->logical_offset_end = tmp_update->logical_offset + tmp_update->size; - tmp_update->flags = c_op->u.write.flags; - if(tmp_update->size >= c_op->u.write.obj_sizes[obj_index]) + tmp_update->flags = c_op->work.u.write.flags; + if(tmp_update->size >= c_op->work.u.write.obj_sizes[obj_index]) { /* mem region is bigger than obj region */ - tmp_update->size = c_op->u.write.obj_sizes[obj_index]; + tmp_update->size = c_op->work.u.write.obj_sizes[obj_index]; } /* advance pointers */ done = advance_listio_ptrs( - c_op->u.write.mem_offsets, - c_op->u.write.mem_sizes, - c_op->u.write.mem_count, + c_op->work.u.write.mem_offsets, + c_op->work.u.write.mem_sizes, + c_op->work.u.write.mem_count, &mem_index, - c_op->u.write.obj_offsets, - c_op->u.write.obj_sizes, - c_op->u.write.obj_count, + c_op->work.u.write.obj_offsets, + c_op->work.u.write.obj_sizes, + c_op->work.u.write.obj_count, &obj_index, tmp_update->size); @@ -1129,7 +1188,7 @@ static int write_op_worker(struct gs_op* op) gs_mutex_lock(&global_log_mutex); if(global_log_offset < 0) { - COSD_INIT_DBT(key, c_op->u.write.oid); + COSD_INIT_DBT(key, c_op->work.u.write.oid); COSD_INIT_DBT(value, global_log_offset); write_op_txn: @@ -1212,7 +1271,7 @@ write_op_txn: /* track this update in the txn accumulator */ gs_mutex_lock(&txn_mutex); - hash_link = gs_hash_search(txn_table, &c_op->u.write.txn_number); + hash_link = gs_hash_search(txn_table, &c_op->work.u.write.txn_number); if(!hash_link) { /* TODO: txn is gone (which could be normal); error handling */ @@ -1223,8 +1282,8 @@ write_op_txn: txn_acc = gs_hash_get_entry(hash_link, struct txn_accumulator, hash_link); /* check that oid and fork match */ - if(txn_acc->oid != c_op->u.write.oid || - txn_acc->fork != c_op->u.write.fork) + if(txn_acc->oid != c_op->work.u.write.oid || + txn_acc->fork != c_op->work.u.write.fork) { /* TODO: graceful error handling; the write operation isn't * legal in this transaction @@ -1390,12 +1449,12 @@ static int create_op_worker(struct gs_op* op) assert(c_op); /* only support one object for now */ - assert(c_op->u.create.requested_oid == 1); - *c_op->u.create.out_oid = 1; + assert(c_op->work.u.create.requested_oid == 1); + *c_op->work.u.create.out_oid = 1; /* create a log file */ sprintf(log_name, "%s/%llu.dat", cosd_log_path, - llu(c_op->u.create.requested_oid)); + llu(c_op->work.u.create.requested_oid)); ret = open(log_name, O_RDWR|O_CREAT|O_DIRECT|O_EXCL|O_NOATIME, S_IRUSR|S_IWUSR); if(ret < 0) @@ -1416,7 +1475,7 @@ static int create_op_worker(struct gs_op* op) /* set initial version of 1 */ val = 1; - COSD_INIT_DBT(key, c_op->u.create.requested_oid); + COSD_INIT_DBT(key, c_op->work.u.create.requested_oid); COSD_INIT_DBT(value, val); ret = ver_dbp->put(ver_dbp, txn, &key, &value, 0); @@ -1447,6 +1506,54 @@ static int create_op_worker(struct gs_op* op) return(1); } +static int cosd_write_set_work( + struct cosd_work* work, + uint64_t oid, + uint64_t fork, + uint64_t txn_number, + char** mem_offsets, + int64_t* mem_sizes, + int mem_count, + int64_t* obj_offsets, + int64_t* obj_sizes, + int obj_count, + int flags) +{ + work->u.write.oid = oid; + work->u.write.fork = fork; + work->u.write.txn_number = txn_number; + work->u.write.mem_count = mem_count; + work->u.write.obj_count = obj_count; + work->u.write.flags = flags; + + /* copy arrays so that we can modify them internally */ + /* TODO: error handling */ + work->u.write.mem_offsets = malloc(mem_count*sizeof(*mem_offsets)); + assert(work->u.write.mem_offsets); + memcpy(work->u.write.mem_offsets, mem_offsets, + mem_count*sizeof(*mem_offsets)); + + work->u.write.mem_sizes = malloc(mem_count*sizeof(*mem_sizes)); + assert(work->u.write.mem_sizes); + memcpy(work->u.write.mem_sizes, mem_sizes, + mem_count*sizeof(*mem_sizes)); + + work->u.write.obj_offsets = malloc(obj_count*sizeof(*obj_offsets)); + assert(work->u.write.obj_offsets); + memcpy(work->u.write.obj_offsets, obj_offsets, + obj_count*sizeof(*obj_offsets)); + + work->u.write.obj_sizes = malloc(obj_count*sizeof(*obj_sizes)); + assert(work->u.write.obj_sizes); + memcpy(work->u.write.obj_sizes, obj_sizes, + obj_count*sizeof(*obj_sizes)); + + work->op_worker = write_op_worker; + work->cleanup_fn = write_op_cleanup; + + return(0); +} + gs_ret_t gs_cosd_write_post( uint64_t oid, uint64_t fork, @@ -1466,47 +1573,64 @@ gs_ret_t gs_cosd_write_post( { struct gs_op *op; struct cosd_op *c_op; + int ret; op = gs_opcache_get(cosd_opcache); gs_op_fill(op, callback, user_ptr, hints, ctx); c_op = gs_op_entry(op, struct cosd_op, op); c_op->op_id = gs_id_gen(gs_cosd_resource_id, (uint64_t)(op->cache_id)); - c_op->u.write.oid = oid; - c_op->u.write.fork = fork; - c_op->u.write.txn_number = txn_number; - c_op->u.write.mem_count = mem_count; - c_op->u.write.obj_count = obj_count; - c_op->u.write.flags = flags; - /* copy arrays so that we can modify them internally */ - /* TODO: error handling */ - c_op->u.write.mem_offsets = malloc(mem_count*sizeof(*mem_offsets)); - assert(c_op->u.write.mem_offsets); - memcpy(c_op->u.write.mem_offsets, mem_offsets, - mem_count*sizeof(*mem_offsets)); + if(flags & COSD_FLAG_AUTO_TXN) + { + /* txn open never blocks; just call it directly */ + ret = gs_cosd_txn_open(oid, fork, txn_number); + if(ret < 0) + { + return(ret); + } - c_op->u.write.mem_sizes = malloc(mem_count*sizeof(*mem_sizes)); - assert(c_op->u.write.mem_sizes); - memcpy(c_op->u.write.mem_sizes, mem_sizes, - mem_count*sizeof(*mem_sizes)); + /* allocate array to describe the write and txn close operations + * that we want to bundle together + */ + c_op->work_array = malloc(sizeof(*c_op->work_array) * 2); + /* TODO: error handling */ + assert(c_op->work_array); + c_op->work_array_count = 2; - c_op->u.write.obj_offsets = malloc(obj_count*sizeof(*obj_offsets)); - assert(c_op->u.write.obj_offsets); - memcpy(c_op->u.write.obj_offsets, obj_offsets, - obj_count*sizeof(*obj_offsets)); + /* set up write operation */ + ret = cosd_write_set_work(&c_op->work_array[0], oid, fork, txn_number, + mem_offsets, mem_sizes, mem_count, obj_offsets, obj_sizes, + obj_count, flags); + if(ret < 0) + { + free(c_op->work_array); + return(ret); + } - c_op->u.write.obj_sizes = malloc(obj_count*sizeof(*obj_sizes)); - assert(c_op->u.write.obj_sizes); - memcpy(c_op->u.write.obj_sizes, obj_sizes, - obj_count*sizeof(*obj_sizes)); + /* set up txn close operation */ + ret = cosd_txn_close_set_work(&c_op->work_array[1], oid, fork, + txn_number); + if(ret < 0) + { + c_op->work_array[0].cleanup_fn(&c_op->work_array[0]); + free(c_op->work_array); + return(ret); + } + } + else + { + ret = cosd_write_set_work(&c_op->work, oid, fork, txn_number, + mem_offsets, mem_sizes, mem_count, obj_offsets, obj_sizes, + obj_count, flags); + if(ret < 0) + { + return(ret); + } + } *op_id = c_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = write_op_worker; - c_op->cleanup_fn = write_op_cleanup; - gs_cosd_launch_op(op, gs_cosd_progress_mode); return 0; @@ -1552,40 +1676,39 @@ gs_ret_t gs_cosd_read_post( c_op = gs_op_entry(op, struct cosd_op, op); c_op->op_id = gs_id_gen(gs_cosd_resource_id, (uint64_t)(op->cache_id)); - c_op->u.read.oid = oid; - c_op->u.read.fork = fork; - c_op->u.read.mem_count = mem_count; - c_op->u.read.obj_count = obj_count; - c_op->u.read.out_size = out_size; + c_op->work.u.read.oid = oid; + c_op->work.u.read.fork = fork; + c_op->work.u.read.mem_count = mem_count; + c_op->work.u.read.obj_count = obj_count; + c_op->work.u.read.out_size = out_size; *out_size = 0; /* copy arrays so that we can modify them internally */ /* TODO: error handling */ - c_op->u.read.mem_offsets = malloc(mem_count*sizeof(*mem_offsets)); - assert(c_op->u.read.mem_offsets); - memcpy(c_op->u.read.mem_offsets, mem_offsets, + c_op->work.u.read.mem_offsets = malloc(mem_count*sizeof(*mem_offsets)); + assert(c_op->work.u.read.mem_offsets); + memcpy(c_op->work.u.read.mem_offsets, mem_offsets, mem_count*sizeof(*mem_offsets)); - c_op->u.read.mem_sizes = malloc(mem_count*sizeof(*mem_sizes)); - assert(c_op->u.read.mem_sizes); - memcpy(c_op->u.read.mem_sizes, mem_sizes, + c_op->work.u.read.mem_sizes = malloc(mem_count*sizeof(*mem_sizes)); + assert(c_op->work.u.read.mem_sizes); + memcpy(c_op->work.u.read.mem_sizes, mem_sizes, mem_count*sizeof(*mem_sizes)); - c_op->u.read.obj_offsets = malloc(obj_count*sizeof(*obj_offsets)); - assert(c_op->u.read.obj_offsets); - memcpy(c_op->u.read.obj_offsets, obj_offsets, + c_op->work.u.read.obj_offsets = malloc(obj_count*sizeof(*obj_offsets)); + assert(c_op->work.u.read.obj_offsets); + memcpy(c_op->work.u.read.obj_offsets, obj_offsets, obj_count*sizeof(*obj_offsets)); - c_op->u.read.obj_sizes = malloc(obj_count*sizeof(*obj_sizes)); - assert(c_op->u.read.obj_sizes); - memcpy(c_op->u.read.obj_sizes, obj_sizes, + c_op->work.u.read.obj_sizes = malloc(obj_count*sizeof(*obj_sizes)); + assert(c_op->work.u.read.obj_sizes); + memcpy(c_op->work.u.read.obj_sizes, obj_sizes, obj_count*sizeof(*obj_sizes)); *op_id = c_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = read_op_worker; - c_op->cleanup_fn = read_op_cleanup; + c_op->work.op_worker = read_op_worker; + c_op->work.cleanup_fn = read_op_cleanup; gs_cosd_launch_op(op, gs_cosd_progress_mode); @@ -1627,9 +1750,8 @@ gs_ret_t gs_cosd_dump_post( *op_id = c_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = dump_op_worker; - c_op->cleanup_fn = NULL; + c_op->work.op_worker = dump_op_worker; + c_op->work.cleanup_fn = NULL; gs_cosd_launch_op(op, gs_cosd_progress_mode); @@ -1660,14 +1782,13 @@ gs_ret_t gs_cosd_create_post( c_op = gs_op_entry(op, struct cosd_op, op); c_op->op_id = gs_id_gen(gs_cosd_resource_id, (uint64_t)(op->cache_id)); - c_op->u.create.requested_oid = requested_oid; - c_op->u.create.out_oid = out_oid; + c_op->work.u.create.requested_oid = requested_oid; + c_op->work.u.create.out_oid = out_oid; *op_id = c_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = create_op_worker; - c_op->cleanup_fn = NULL; + c_op->work.op_worker = create_op_worker; + c_op->work.cleanup_fn = NULL; gs_cosd_launch_op(op, gs_cosd_progress_mode); @@ -1724,7 +1845,7 @@ static int get_version_op_worker(struct gs_op* op) assert(c_op); /* only support one object for now */ - assert(c_op->u.get_version.oid == 1); + assert(c_op->work.u.get_version.oid == 1); /* do db stuff */ ret = envp->txn_begin(envp, NULL, &txn, 0); @@ -1736,7 +1857,7 @@ static int get_version_op_worker(struct gs_op* op) } /* read current version */ - COSD_INIT_DBT(key, c_op->u.get_version.oid); + COSD_INIT_DBT(key, c_op->work.u.get_version.oid); COSD_INIT_DBT(value, version); ret = ver_dbp->get(ver_dbp, txn, &key, &value, 0); @@ -1746,7 +1867,7 @@ static int get_version_op_worker(struct gs_op* op) assert(0); } - *c_op->u.get_version.version = version; + *c_op->work.u.get_version.version = version; ret = txn->commit(txn, 0); if(ret != 0) @@ -1778,14 +1899,13 @@ gs_ret_t gs_cosd_get_version_post( c_op = gs_op_entry(op, struct cosd_op, op); c_op->op_id = gs_id_gen(gs_cosd_resource_id, (uint64_t)(op->cache_id)); - c_op->u.get_version.oid = oid; - c_op->u.get_version.version = version; + c_op->work.u.get_version.oid = oid; + c_op->work.u.get_version.version = version; *op_id = c_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = get_version_op_worker; - c_op->cleanup_fn = NULL; + c_op->work.op_worker = get_version_op_worker; + c_op->work.cleanup_fn = NULL; gs_cosd_launch_op(op, gs_cosd_progress_mode); @@ -1835,14 +1955,14 @@ static int txn_close_op_worker(struct gs_op* op) assert(c_op); /* only support one object for now */ - assert(c_op->u.txn_close.oid == 1); + assert(c_op->work.u.txn_close.oid == 1); - lmk.oid = c_op->u.txn_close.oid; - lmk.fork = c_op->u.txn_close.fork; + lmk.oid = c_op->work.u.txn_close.oid; + lmk.fork = c_op->work.u.txn_close.fork; /* pull txn accumulator out of hash so no one can touch it */ gs_mutex_lock(&txn_mutex); - hash_link = gs_hash_search(txn_table, &c_op->u.txn_close.txn_number); + hash_link = gs_hash_search(txn_table, &c_op->work.u.txn_close.txn_number); if(!hash_link) { /* TODO: txn is gone (which could be normal); error handling */ @@ -1853,8 +1973,8 @@ static int txn_close_op_worker(struct gs_op* op) txn_acc = gs_hash_get_entry(hash_link, struct txn_accumulator, hash_link); - if(txn_acc->oid != c_op->u.txn_close.oid || - txn_acc->fork != c_op->u.txn_close.fork) + if(txn_acc->oid != c_op->work.u.txn_close.oid || + txn_acc->fork != c_op->work.u.txn_close.fork) { /* TODO: error handling */ /* mismatch between txn open and close */ @@ -1913,8 +2033,8 @@ txn_close_op_retry: tmp_update = gs_list_get_entry(iterator, struct logical_map_entry, list_link); - lmk.oid = c_op->u.txn_close.oid; - lmk.fork = c_op->u.txn_close.fork; + lmk.oid = c_op->work.u.txn_close.oid; + lmk.fork = c_op->work.u.txn_close.fork; lmk.logical_offset_end = tmp_update->logical_offset + 1; done = 0; c_get_flag = DB_SET_RANGE; @@ -1947,8 +2067,8 @@ txn_close_op_retry: * branch, or if we go past our offset. Ignore the offset if * flag is set to truncate the whole fork. */ - if(lmk.oid != c_op->u.txn_close.oid || - lmk.fork != c_op->u.txn_close.fork || + if(lmk.oid != c_op->work.u.txn_close.oid || + lmk.fork != c_op->work.u.txn_close.fork || (lme.logical_offset >= tmp_update->logical_offset_end && !(tmp_update->flags & COSD_FLAG_TRUNC_WRITE))) { @@ -2022,8 +2142,8 @@ txn_close_op_retry: tmp_entry = gs_list_get_entry(iterator, struct logical_map_entry, list_link); lme = *tmp_entry; - lmk.oid = c_op->u.txn_close.oid; - lmk.fork = c_op->u.txn_close.fork; + lmk.oid = c_op->work.u.txn_close.oid; + lmk.fork = c_op->work.u.txn_close.fork; lmk.logical_offset_end = lme.logical_offset_end; ret = log_map_dbp->put(log_map_dbp, txn, &key, &value, 0); if(ret != 0) @@ -2042,7 +2162,7 @@ txn_close_op_retry: free(free_ptr); /* read current version */ - COSD_INIT_DBT(key, c_op->u.txn_close.oid); + COSD_INIT_DBT(key, c_op->work.u.txn_close.oid); COSD_INIT_DBT(value, version); ret = ver_dbp->get(ver_dbp, txn, &key, &value, DB_RMW); @@ -2069,7 +2189,7 @@ txn_close_op_retry: missing_ver = version + 1; while(missing_ver < txn_acc->txn_number) { - mv.oid = c_op->u.txn_close.oid; + mv.oid = c_op->work.u.txn_close.oid; mv.version = missing_ver; /* record each missing ver in db */ @@ -2111,7 +2231,7 @@ txn_close_op_retry: memset(&mv_key, 0, sizeof(DBT)); mv_key.data = &mv; mv_key.size = sizeof(mv); - mv.oid = c_op->u.txn_close.oid; + mv.oid = c_op->work.u.txn_close.oid; mv.version = txn_acc->txn_number; /* delete from db db */ @@ -2135,7 +2255,7 @@ txn_close_op_retry: * relative to txn? */ gs_mutex_lock(&global_log_mutex); - COSD_INIT_DBT(key, c_op->u.txn_close.oid); + COSD_INIT_DBT(key, c_op->work.u.txn_close.oid); COSD_INIT_DBT(value, global_log_offset); ret = log_offset_dbp->put(log_offset_dbp, txn, &key, &value, 0); @@ -2187,6 +2307,22 @@ txn_close_op_retry: return(1); } +static int cosd_txn_close_set_work( + struct cosd_work* work, + uint64_t oid, + uint64_t fork, + uint64_t txn_number) +{ + work->u.txn_close.oid = oid; + work->u.txn_close.fork = fork; + work->u.txn_close.txn_number = txn_number; + + work->op_worker = txn_close_op_worker; + work->cleanup_fn = NULL; + + return(0); +} + gs_ret_t gs_cosd_txn_close_post( uint64_t oid, uint64_t fork, @@ -2199,22 +2335,19 @@ gs_ret_t gs_cosd_txn_close_post( { struct gs_op *op; struct cosd_op *c_op; + int ret; op = gs_opcache_get(cosd_opcache); gs_op_fill(op, callback, user_ptr, hints, ctx); c_op = gs_op_entry(op, struct cosd_op, op); c_op->op_id = gs_id_gen(gs_cosd_resource_id, (uint64_t)(op->cache_id)); - c_op->u.txn_close.oid = oid; - c_op->u.txn_close.fork = fork; - c_op->u.txn_close.txn_number = txn_number; - - *op_id = c_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = txn_close_op_worker; - c_op->cleanup_fn = NULL; + ret = cosd_txn_close_set_work(&c_op->work, oid, fork, txn_number); + if(ret < 0) + return(ret); + *op_id = c_op->op_id; gs_cosd_launch_op(op, gs_cosd_progress_mode); return 0; diff --git a/code/src/gsl/resources/cosd-prototype/cosd-prototype.gsh b/code/src/gsl/resources/cosd-prototype/cosd-prototype.gsh index 036a2f5..a16ba06 100644 --- a/code/src/gsl/resources/cosd-prototype/cosd-prototype.gsh +++ b/code/src/gsl/resources/cosd-prototype/cosd-prototype.gsh @@ -20,6 +20,7 @@ #include "include/gs.h" #define COSD_FLAG_TRUNC_WRITE 1 /**< truncate fork on write */ +#define COSD_FLAG_AUTO_TXN 2 /**< automatically wrap write in a transaction */ /** modes of making progress on posted storage operations */ enum progress_mode diff --git a/code/src/gsl/resources/cosd-prototype/test/cosd1.gs b/code/src/gsl/resources/cosd-prototype/test/cosd1.gs index 39c782e..c4a466f 100644 --- a/code/src/gsl/resources/cosd-prototype/test/cosd1.gs +++ b/code/src/gsl/resources/cosd-prototype/test/cosd1.gs @@ -401,6 +401,27 @@ static __blocking int do_cosd_test(void) ret = gs_cosd_dump(); assert(ret == 0); + /********************************************************/ + + printf("writing 512-1024 with auto txn...\n"); + buffer_offsets[0] = buffer; + buffer_szs[0] = 512; + obj_offsets[0] = 512; + obj_szs[0] = 512; + ret = gs_cosd_write(1, 0, (++version), buffer_offsets, buffer_szs, 1, + obj_offsets, obj_szs, 1, COSD_FLAG_AUTO_TXN); + if(ret != 0) + { + printf("Error writing 512-1024: %d\n", ret); + free(buffer); + return 1; + } + printf("DONE\n"); + + ret = gs_cosd_dump(); + assert(ret == 0); + + free(buffer); return 0; } diff --git a/code/src/gsl/resources/storage/gs-storage.c b/code/src/gsl/resources/storage/gs-storage.c index 74bea89..0ce1509 100644 --- a/code/src/gsl/resources/storage/gs-storage.c +++ b/code/src/gsl/resources/storage/gs-storage.c @@ -53,6 +53,7 @@ struct storage_op } getattr_list; } u; + int (*op_worker)(struct gs_op* op); gs_op_id_t op_id; struct gs_op op; int error_code; @@ -164,7 +165,7 @@ static void* thread_fn(void* foo) gs_mutex_unlock(&storage_mutex); /* spin on op worker */ - while(op->op_worker(op) != 1); + while(sto_op->op_worker(op) != 1); /* done */ gs_opcache_complete_op(storage_opcache, op, sto_op->error_code); @@ -242,12 +243,10 @@ gs_ret_t gs_storage_remove_post( sto_op = gs_op_entry(op, struct storage_op, op); sto_op->op_id = gs_id_gen(gs_storage_resource_id, (uint64_t)(op->cache_id)); sto_op->u.remove.oid = oid; + sto_op->op_worker = remove_op_worker; *op_id = sto_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = remove_op_worker; - gs_storage_launch_op(op, gs_storage_progress_mode); return 0; @@ -280,12 +279,10 @@ gs_ret_t gs_storage_create_post( sto_op->op_id = gs_id_gen(gs_storage_resource_id, (uint64_t)(op->cache_id)); sto_op->u.create.requested_oid = requested_oid; sto_op->u.create.out_oid = out_oid; + sto_op->op_worker = create_op_worker; *op_id = sto_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = create_op_worker; - gs_storage_launch_op(op, gs_storage_progress_mode); return 0; @@ -326,12 +323,10 @@ gs_ret_t gs_storage_setattr_list_post( sto_op->u.setattr_list.number_array = number_array; sto_op->u.setattr_list.val_array = val_array; sto_op->u.setattr_list.len_array = len_array; + sto_op->op_worker = setattr_list_op_worker; *op_id = sto_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = setattr_list_op_worker; - gs_storage_launch_op(op, gs_storage_progress_mode); return(0); @@ -376,12 +371,10 @@ gs_ret_t gs_storage_getattr_list_post( sto_op->u.getattr_list.number_array = number_array; sto_op->u.getattr_list.val_array = val_array; sto_op->u.getattr_list.len_array = len_array; + sto_op->op_worker = getattr_list_op_worker; *op_id = sto_op->op_id; - /* TODO: put this in the fill function if we keep it? */ - op->op_worker = getattr_list_op_worker; - gs_storage_launch_op(op, gs_storage_progress_mode); return(0); @@ -435,7 +428,7 @@ static int gs_storage_poll(gs_context_t context, int millisecs) sto_op = gs_op_entry(gop, struct storage_op, op); - ret = gop->op_worker(gop); + ret = sto_op->op_worker(gop); if(ret == 1) { /* done */ hooks/post-receive -- Grayskull Repository
participants (1)
-
noreply@mcs.anl.gov