diff --git a/c/Makefile b/c/Makefile index b7fe98c5c..135b47787 100644 --- a/c/Makefile +++ b/c/Makefile @@ -21,7 +21,10 @@ input/lrec_reader_mmap_nidx.c input/lrec_reader_stdio_nidx.c \ input/lrec_reader_mmap_xtab.c input/lrec_reader_stdio_xtab.c \ containers/test_lrec.c -TEST_JOIN_BUCKET_KEEPER_SRCS = lib/mlrutil.c lib/mlr_globals.c containers/lrec.c \ +TEST_JOIN_BUCKET_KEEPER_SRCS = \ +lib/mlrutil.c lib/mlr_globals.c \ +mapping/context.c \ +containers/lrec.c \ containers/sllv.c containers/slls.c containers/lhmslv.c containers/hss.c containers/mixutil.c \ containers/header_keeper.c \ containers/join_bucket_keeper.c \ @@ -88,6 +91,9 @@ unit-test: test-mlrutil test-lrec test-join-bucket-keeper reg-test: ./test/run +tj: test-join-bucket-keeper + ./test-join-bucket-keeper + # Run this after unit-test expected output has changed, and is verified to be # OK. (Example: after adding new test cases in test/run.) regtest-copy: diff --git a/c/containers/join_bucket_keeper.c b/c/containers/join_bucket_keeper.c index e05dbaef4..adba4cc96 100644 --- a/c/containers/join_bucket_keeper.c +++ b/c/containers/join_bucket_keeper.c @@ -1,6 +1,7 @@ #include #include "lib/mlrutil.h" #include "lib/mlr_globals.h" +#include "mapping/context.h" #include "containers/mixutil.h" #include "containers/join_bucket_keeper.h" #include "input/lrec_readers.h" @@ -18,6 +19,10 @@ // (3) leof: Lv == null, peek == null, leof = true // ---------------------------------------------------------------- +// Private methods +static int join_bucket_keeper_get_state(join_bucket_keeper_t* pkeeper); +static void join_bucket_keeper_initial_fill(join_bucket_keeper_t* pkeeper); + // ---------------------------------------------------------------- join_bucket_keeper_t* join_bucket_keeper_alloc( char* left_file_name, @@ -31,19 +36,31 @@ join_bucket_keeper_t* join_bucket_keeper_alloc( slls_t* pleft_field_names ) { - join_bucket_keeper_t* pkeeper = mlr_malloc_or_die(sizeof(join_bucket_keeper_t)); - - pkeeper->plrec_reader = lrec_reader_alloc(input_file_format, use_mmap_for_read, + lrec_reader_t* plrec_reader = lrec_reader_alloc(input_file_format, use_mmap_for_read, irs, ifs, allow_repeat_ifs, ips, allow_repeat_ips); - pkeeper->pvhandle = pkeeper->plrec_reader->popen_func(left_file_name); - pkeeper->plrec_reader->psof_func(pkeeper->plrec_reader->pvstate); + void* pvhandle = plrec_reader->popen_func(left_file_name); // xxx move this ... + plrec_reader->psof_func(plrec_reader->pvstate); // xxx move this ... - pkeeper->pctx = mlr_malloc_or_die(sizeof(context_t)); - pkeeper->pctx->nr = 0; // xxx make an init func & use it here & in stream.c? - pkeeper->pctx->fnr = 0; // xxx incr this in the readers ... - pkeeper->pctx->filenum = 1; - pkeeper->pctx->filename = left_file_name; + context_t* pctx = mlr_malloc_or_die(sizeof(context_t)); + context_init(pctx, left_file_name); + + return join_bucket_keeper_alloc_from_reader(plrec_reader, pvhandle, pctx, pleft_field_names); +} + +// ---------------------------------------------------------------- +join_bucket_keeper_t* join_bucket_keeper_alloc_from_reader( + lrec_reader_t* plrec_reader, + void* pvhandle, + context_t* pctx, + slls_t* pleft_field_names +) { + + join_bucket_keeper_t* pkeeper = mlr_malloc_or_die(sizeof(join_bucket_keeper_t)); + + pkeeper->plrec_reader = plrec_reader; + pkeeper->pvhandle = pvhandle; + pkeeper->pctx = pctx; pkeeper->pleft_field_names = slls_copy(pleft_field_names); // xxx be sure the caller frees its own pkeeper->pleft_field_values = NULL; @@ -65,6 +82,104 @@ void join_bucket_keeper_free(join_bucket_keeper_t* pkeeper) { pkeeper->plrec_reader->pclose_func(pkeeper->pvhandle); } +// ---------------------------------------------------------------- +// xxx cmt re who frees +void join_bucket_keeper_emit(join_bucket_keeper_t* pkeeper, slls_t* pright_field_values, + sllv_t** ppbucket_paired, sllv_t** ppbucket_left_unpaired) +{ + *ppbucket_paired = NULL; + *ppbucket_left_unpaired = NULL; + int cmp = 0; + + if (pkeeper->state == LEFT_STATE_0_PREFILL) { + // try fill Lv & peek; next state is 1,2,3 & continue from there. + join_bucket_keeper_initial_fill(pkeeper); + pkeeper->state = join_bucket_keeper_get_state(pkeeper); + } + + // Return the final left-unpaireds after right EOF. + if (pright_field_values == NULL) { + *ppbucket_left_unpaired = pkeeper->precords; + return; + } + + switch (pkeeper->state) { + case LEFT_STATE_1_FULL: + case LEFT_STATE_2_LAST_BUCKET: // Intentional fall-through + cmp = slls_compare_lexically(pkeeper->pleft_field_values, pright_field_values); + if (cmp < 0) { + *ppbucket_left_unpaired = pkeeper->precords; + // xxx advance left ... + } else if (cmp == 0) { + *ppbucket_paired = pkeeper->precords; + } + + break; + + case LEFT_STATE_3_EOF: + break; + + default: + fprintf(stderr, "%s: internal coding error: failed transition from prefill state.\n", + MLR_GLOBALS.argv0); + exit(1); + break; + } + + pkeeper->state = join_bucket_keeper_get_state(pkeeper); + +} + +// // xxx cmt why rec-peek +// if (!pkeeper->leof && pkeeper->prec_peek == NULL) { +// pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, +// pkeeper->plrec_reader->pvstate, pkeeper->pctx); +// if (pkeeper->prec_peek == NULL) +// pkeeper->leof = TRUE; +// } +// +// // xxx not quite: only if *already* leof. +// if (pkeeper->leof) { +// int cmp = slls_compare_lexically(pkeeper->pbucket->pjoin_values, pright_field_values); +// // xxx think through various cases +// if (/* xxx stub */ cmp == 999) { +// // rename: "bucket" used at different nesting levels. it's confusing. +// *ppbucket_paired = pkeeper->precords; +// } else { +// *ppbucket_left_unpaired = pkeeper->precords; +// } +// pkeeper->pbucket = NULL; +// return; +// } +// +// // xxx now that we've got a peek record: fill the bucket with like records +// // until there's a non-like on-deck. +// // +// // xxx do this only if it's time for a change. +// // +// // xxx rename "peek" to "on-deck"? +// +//#if 0 +// sllv_empty(pkeeper->pbucket); +// while (TRUE) { +// sllv_add(pkeeper->pbucket, pkeeper->prec_peek); +// pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, +// pkeeper->plrec_reader->pvstate, pkeeper->pctx); +// get selected keys +// if (pkeeper->prec_peek == NULL) +// pkeeper->leof = TRUE; +// x +// } +//#endif +// +// // xxx stub +// lrec_t* pleft_rec = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, pkeeper->plrec_reader->pvstate, +// pkeeper->pctx); +// sllv_t* pfoo = sllv_alloc(); +// if (pleft_rec != NULL) +// sllv_add(pfoo, pleft_rec); +// *ppbucket_paired = pfoo; + // ---------------------------------------------------------------- // xxx left state: // (0) pre-fill: Lv == null, peek == null, leof = false @@ -159,100 +274,6 @@ static void join_bucket_keeper_initial_fill(join_bucket_keeper_t* pkeeper) { // sllv_add(pfoo, pleft_rec); // *ppbucket_paired = pfoo; -// xxx cmt re who frees -void join_bucket_keeper_emit(join_bucket_keeper_t* pkeeper, slls_t* pright_field_values, - sllv_t** ppbucket_paired, sllv_t** ppbucket_left_unpaired) -{ - *ppbucket_paired = NULL; - *ppbucket_left_unpaired = NULL; - int cmp = 0; - - if (pkeeper->state == LEFT_STATE_0_PREFILL) { - // try fill Lv & peek; next state is 1,2,3 & continue from there. - join_bucket_keeper_initial_fill(pkeeper); - pkeeper->state = join_bucket_keeper_get_state(pkeeper); - } - - // xxx drain on pright_field_values == NULL, for returning the final - // left-unpaireds after right EOF. - - switch (pkeeper->state) { - case LEFT_STATE_1_FULL: - case LEFT_STATE_2_LAST_BUCKET: // Intentional fall-through - cmp = slls_compare_lexically(pkeeper->pleft_field_values, pright_field_values); - if (cmp < 0) { - *ppbucket_left_unpaired = pkeeper->precords; - // xxx advance left ... - } else if (cmp == 0) { - *ppbucket_paired = pkeeper->precords; - } - - break; - - case LEFT_STATE_3_EOF: - break; - - default: - fprintf(stderr, "%s: internal coding error: failed transition from prefill state.\n", - MLR_GLOBALS.argv0); - exit(1); - break; - } - - pkeeper->state = join_bucket_keeper_get_state(pkeeper); - -// // xxx cmt why rec-peek -// if (!pkeeper->leof && pkeeper->prec_peek == NULL) { -// pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, -// pkeeper->plrec_reader->pvstate, pkeeper->pctx); -// if (pkeeper->prec_peek == NULL) -// pkeeper->leof = TRUE; -// } -// -// // xxx not quite: only if *already* leof. -// if (pkeeper->leof) { -// int cmp = slls_compare_lexically(pkeeper->pbucket->pjoin_values, pright_field_values); -// // xxx think through various cases -// if (/* xxx stub */ cmp == 999) { -// // rename: "bucket" used at different nesting levels. it's confusing. -// *ppbucket_paired = pkeeper->precords; -// } else { -// *ppbucket_left_unpaired = pkeeper->precords; -// } -// pkeeper->pbucket = NULL; -// return; -// } -// -// // xxx now that we've got a peek record: fill the bucket with like records -// // until there's a non-like on-deck. -// // -// // xxx do this only if it's time for a change. -// // -// // xxx rename "peek" to "on-deck"? -// -//#if 0 -// sllv_empty(pkeeper->pbucket); -// while (TRUE) { -// sllv_add(pkeeper->pbucket, pkeeper->prec_peek); -// pkeeper->prec_peek = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, -// pkeeper->plrec_reader->pvstate, pkeeper->pctx); -// get selected keys -// if (pkeeper->prec_peek == NULL) -// pkeeper->leof = TRUE; -// x -// } -//#endif -// -// // xxx stub -// lrec_t* pleft_rec = pkeeper->plrec_reader->pprocess_func(pkeeper->pvhandle, pkeeper->plrec_reader->pvstate, -// pkeeper->pctx); -// sllv_t* pfoo = sllv_alloc(); -// if (pleft_rec != NULL) -// sllv_add(pfoo, pleft_rec); -// *ppbucket_paired = pfoo; - -} - // ---------------------------------------------------------------- // +-----------+-----------+-----------+-----------+-----------+-----------+ diff --git a/c/containers/join_bucket_keeper.h b/c/containers/join_bucket_keeper.h index 5ee55fdda..d9520e7dc 100644 --- a/c/containers/join_bucket_keeper.h +++ b/c/containers/join_bucket_keeper.h @@ -13,12 +13,11 @@ typedef struct _join_bucket_keeper_t { context_t* pctx; slls_t* pleft_field_names; - - int state; slls_t* pleft_field_values; sllv_t* precords; lrec_t* prec_peek; int leof; + int state; } join_bucket_keeper_t; @@ -31,12 +30,20 @@ join_bucket_keeper_t* join_bucket_keeper_alloc( int allow_repeat_ifs, char ips, int allow_repeat_ips, - slls_t* pleft_field_names -); + slls_t* pleft_field_names); + +join_bucket_keeper_t* join_bucket_keeper_alloc_from_reader( + lrec_reader_t* plrec_reader, + void* pvhandle, + context_t* pctx, + slls_t* pleft_field_names); void join_bucket_keeper_free(join_bucket_keeper_t* pkeeper); -void join_bucket_keeper_emit(join_bucket_keeper_t* pkeeper, slls_t* pright_field_values, - sllv_t** ppbucket_paired, sllv_t** ppbucket_left_unpaired); +void join_bucket_keeper_emit( + join_bucket_keeper_t* pkeeper, + slls_t* pright_field_values, + sllv_t** ppbucket_paired, + sllv_t** ppbucket_left_unpaired); #endif // JOIN_BUCKET_KEEPER_H diff --git a/c/containers/test_join_bucket_keeper.c b/c/containers/test_join_bucket_keeper.c index 80b6b31b0..93fcf63b3 100644 --- a/c/containers/test_join_bucket_keeper.c +++ b/c/containers/test_join_bucket_keeper.c @@ -5,6 +5,7 @@ #include "containers/lrec.h" #include "containers/sllv.h" #include "input/lrec_readers.h" +#include "containers/join_bucket_keeper.h" #ifdef __TEST_JOIN_BUCKET_KEEPER_MAIN__ int tests_run = 0; @@ -13,32 +14,97 @@ int assertions_run = 0; int assertions_failed = 0; // ---------------------------------------------------------------- -static char* test_foo() { - sllv_t* precords = sllv_alloc(); - sllv_add(precords, lrec_literal_2("a","1", "b","10")); - sllv_add(precords, lrec_literal_2("a","1", "b","11")); - sllv_add(precords, lrec_literal_2("a","2", "b","12")); - sllv_add(precords, lrec_literal_2("a","2", "b","13")); - sllv_add(precords, lrec_literal_2("a","3", "b","14")); - sllv_add(precords, lrec_literal_2("a","3", "b","15")); - lrec_reader_t* preader = lrec_reader_in_memory_alloc(precords); - printf("#=%d\n", precords->length); - mu_assert_lf(precords->length == 6); - while (TRUE) { - lrec_t* precord = preader->pprocess_func(NULL, preader->pvstate, NULL); - if (precord == NULL) - break; - lrec_print(precord); - } - printf("#=%d\n", precords->length); - mu_assert_lf(precords->length == 0); +static void set_up( + slls_t** ppleft_field_names, + lrec_reader_t** ppreader) +{ + slls_t* pleft_field_names = slls_alloc(); + slls_add_no_free(pleft_field_names, "l"); + sllv_t* precords = sllv_alloc(); + sllv_add(precords, lrec_literal_2("l","1", "b","10")); + sllv_add(precords, lrec_literal_2("l","1", "b","11")); + sllv_add(precords, lrec_literal_2("l","3", "b","12")); + sllv_add(precords, lrec_literal_2("l","3", "b","13")); + sllv_add(precords, lrec_literal_2("l","3", "b","14")); + sllv_add(precords, lrec_literal_2("l","5", "b","15")); + + lrec_reader_t* preader = lrec_reader_in_memory_alloc(precords); + + *ppleft_field_names = pleft_field_names; + *ppreader = preader; +} + +// ---------------------------------------------------------------- +static char* test_bar0() { + printf("test_bar enter\n"); + + slls_t* pleft_field_names; + lrec_reader_t* preader; + set_up(&pleft_field_names, &preader); + + void* pvhandle = NULL; // xxx move these into the jbk obj? + context_t* pctx = NULL; // xxx revisit + + join_bucket_keeper_t* pkeeper = join_bucket_keeper_alloc_from_reader(preader, pvhandle, pctx, + pleft_field_names); + + sllv_t* pbucket_paired; + sllv_t* pbucket_left_unpaired; + + slls_t* pright_field_values = slls_alloc(); + slls_add_no_free(pright_field_values, "0"); + join_bucket_keeper_emit(pkeeper, pright_field_values, &pbucket_paired, &pbucket_left_unpaired); + mu_assert_lf(pbucket_paired == NULL); + mu_assert_lf(pbucket_left_unpaired == NULL); + + pright_field_values = slls_alloc(); + slls_add_no_free(pright_field_values, "1"); + + join_bucket_keeper_emit(pkeeper, pright_field_values, &pbucket_paired, &pbucket_left_unpaired); + mu_assert_lf(pbucket_paired != NULL); + mu_assert_lf(pbucket_paired->length == 2); + mu_assert_lf(pbucket_left_unpaired == NULL); + + printf("test_bar exit\n"); + printf("\n"); + return 0; +} + +// ---------------------------------------------------------------- +static char* test_bar() { + printf("test_bar enter\n"); + + slls_t* pleft_field_names; + lrec_reader_t* preader; + set_up(&pleft_field_names, &preader); + + void* pvhandle = NULL; // xxx move these into the jbk obj? + context_t* pctx = NULL; // xxx revisit + + join_bucket_keeper_t* pkeeper = join_bucket_keeper_alloc_from_reader(preader, pvhandle, pctx, + pleft_field_names); + + sllv_t* pbucket_paired; + sllv_t* pbucket_left_unpaired; + + slls_t* pright_field_values = slls_alloc(); + slls_add_no_free(pright_field_values, "2"); + + join_bucket_keeper_emit(pkeeper, pright_field_values, &pbucket_paired, &pbucket_left_unpaired); + mu_assert_lf(pbucket_paired == NULL); + mu_assert_lf(pbucket_left_unpaired != NULL); + mu_assert_lf(pbucket_left_unpaired->length == 2); + + printf("test_bar exit\n"); + printf("\n"); return 0; } // ================================================================ static char * run_all_tests() { - mu_run_test(test_fotest_foo); + mu_run_test(test_bar0); + mu_run_test(test_bar); return 0; } diff --git a/c/mapping/context.c b/c/mapping/context.c new file mode 100644 index 000000000..9bd820b94 --- /dev/null +++ b/c/mapping/context.c @@ -0,0 +1,8 @@ +#include "context.h" + +void context_init(context_t* pctx, char* first_file_name) { + pctx->nr = 0; + pctx->fnr = 0; + pctx->filenum = 1; + pctx->filename = first_file_name; +} diff --git a/c/mapping/context.h b/c/mapping/context.h index 8f952a932..cd79aecff 100644 --- a/c/mapping/context.h +++ b/c/mapping/context.h @@ -9,4 +9,6 @@ typedef struct _context_t { char* filename; } context_t; +void context_init(context_t* pctx, char* first_file_name); + #endif // CONTEXT_H diff --git a/c/stream/stream.c b/c/stream/stream.c index 1dcadb086..2cfcd361e 100644 --- a/c/stream/stream.c +++ b/c/stream/stream.c @@ -25,16 +25,16 @@ int do_stream_chained(char** filenames, lrec_reader_t* plrec_reader, sllv_t* pma { FILE* output_stream = stdout; - context_t ctx = { .nr = 0, .fnr = 0, .filenum = 0, .filename = NULL }; + context_t ctx = { .nr = 0, .fnr = 0, .filenum = 0, .filename = NULL }; // xxx make a method int ok = 1; if (*filenames == NULL) { - ctx.filenum++; + ctx.filenum++; // xxx make a method ctx.filename = "(stdin)"; ctx.fnr = 0; ok = do_file_chained("-", &ctx, plrec_reader, pmapper_list, plrec_writer, output_stream) && ok; } else { for (char** pfilename = filenames; *pfilename != NULL; pfilename++) { - ctx.filenum++; + ctx.filenum++; // xxx make a method ctx.filename = *pfilename; ctx.fnr = 0; // Start-of-file hook, e.g. expecting CSV headers on input. @@ -63,8 +63,8 @@ static int do_file_chained(char* filename, context_t* pctx, lrec_t* pinrec = plrec_reader->pprocess_func(pvhandle, plrec_reader->pvstate, pctx); if (pinrec == NULL) break; - // incr inside the readers - pctx->nr++; + // xxx incr inside the readers + pctx->nr++; // xxx make a method pctx->fnr++; drive_lrec(pinrec, pctx, pmapper_list->phead, plrec_writer, output_stream); }