mirror of
https://github.com/wassname/ray.git
synced 2026-08-15 12:45:23 +08:00
task queue tests and extensions (#18)
* task queue tests and extensions * clean up test
This commit is contained in:
committed by
Robert Nishihara
parent
313241e303
commit
7a079547b0
@@ -2,6 +2,8 @@
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
|
||||
#include "utarray.h"
|
||||
|
||||
#include "task.h"
|
||||
#include "common.h"
|
||||
#include "io.h"
|
||||
@@ -30,7 +32,8 @@ typedef struct {
|
||||
} task_arg;
|
||||
|
||||
struct task_spec_impl {
|
||||
function_id func_id;
|
||||
/* Function ID of the task. */
|
||||
function_id function_id;
|
||||
/* Total number of arguments. */
|
||||
int64_t num_args;
|
||||
/* Index of the last argument that has been constructed. */
|
||||
@@ -52,14 +55,14 @@ struct task_spec_impl {
|
||||
(sizeof(task_spec) + ((NUM_ARGS) + (NUM_RETURNS)) * sizeof(task_arg) + \
|
||||
(ARGS_VALUE_SIZE))
|
||||
|
||||
task_spec *alloc_task_spec(function_id func_id,
|
||||
task_spec *alloc_task_spec(function_id function_id,
|
||||
int64_t num_args,
|
||||
int64_t num_returns,
|
||||
int64_t args_value_size) {
|
||||
int64_t size = TASK_SPEC_SIZE(num_args, num_returns, args_value_size);
|
||||
task_spec *task = malloc(size);
|
||||
memset(task, 0, size);
|
||||
task->func_id = func_id;
|
||||
task->function_id = function_id;
|
||||
task->num_args = num_args;
|
||||
task->arg_index = 0;
|
||||
task->num_returns = num_returns;
|
||||
@@ -72,6 +75,10 @@ int64_t task_size(task_spec *spec) {
|
||||
spec->args_value_size);
|
||||
}
|
||||
|
||||
unique_id *task_function(task_spec *spec) {
|
||||
return &spec->function_id;
|
||||
}
|
||||
|
||||
int64_t task_num_args(task_spec *spec) {
|
||||
return spec->num_args;
|
||||
}
|
||||
@@ -153,3 +160,86 @@ task_spec *read_task(int fd) {
|
||||
CHECK(task_size(spec) == length);
|
||||
return spec;
|
||||
}
|
||||
|
||||
void print_task(task_spec *spec, UT_string *output) {
|
||||
/* For converting an id to hex, which has double the number
|
||||
* of bytes compared to the id (+ 1 byte for '\0'). */
|
||||
static char hex[2 * UNIQUE_ID_SIZE + 1];
|
||||
/* Print function id. */
|
||||
sha1_to_hex(&task_function(spec)->id[0], &hex[0]);
|
||||
utstring_printf(output, "fun %s ", &hex[0]);
|
||||
/* Print arguments. */
|
||||
for (int i = 0; i < task_num_args(spec); ++i) {
|
||||
sha1_to_hex(&task_arg_id(spec, i)->id[0], &hex[0]);
|
||||
utstring_printf(output, " id:%d %s", i, &hex[0]);
|
||||
}
|
||||
/* Print return ids. */
|
||||
for (int i = 0; i < task_num_returns(spec); ++i) {
|
||||
object_id *object_id = task_return(spec, i);
|
||||
sha1_to_hex(&object_id->id[0], &hex[0]);
|
||||
utstring_printf(output, " ret:%d %s", i, &hex[0]);
|
||||
}
|
||||
}
|
||||
|
||||
UT_icd unique_id_icd = {sizeof(unique_id), NULL, NULL, NULL};
|
||||
|
||||
task_spec *parse_task(char *task_string, int64_t task_length) {
|
||||
/* We make one pass through task_string to store all the argument ids
|
||||
* in "args" and all the return ids in "returns". */
|
||||
UT_array *args;
|
||||
utarray_new(args, &unique_id_icd);
|
||||
UT_array *returns;
|
||||
utarray_new(returns, &unique_id_icd);
|
||||
function_id function_id;
|
||||
char *cursor = strtok(task_string, " ");
|
||||
int index = 0;
|
||||
while (cursor != NULL) {
|
||||
/* This will be equal to "args" or "returns" depending on whether we
|
||||
* are processing an argument id or a return id. */
|
||||
UT_array *target = NULL;
|
||||
if (strncmp("fun", cursor, 3) == 0) {
|
||||
/* Parse function id. */
|
||||
CHECK(cursor + 2 * UNIQUE_ID_SIZE + 1 <= task_string + task_length);
|
||||
cursor = strtok(NULL, " ");
|
||||
hex_to_sha1(cursor, &function_id.id[0]);
|
||||
cursor = strtok(NULL, " ");
|
||||
CHECK(cursor);
|
||||
continue;
|
||||
} else if (strncmp("id:", cursor, 3) == 0) {
|
||||
/* Parse pass by reference argument. */
|
||||
sscanf(cursor, "id:%d", &index);
|
||||
target = args;
|
||||
} else if (strncmp("val:", cursor, 4) == 0) {
|
||||
/* Parse pass by value argument. */
|
||||
sscanf(cursor, "val:%d", &index);
|
||||
CHECK(0); /* Not implemented yet */
|
||||
} else if (strncmp("ret:", cursor, 4) == 0) {
|
||||
/* Parse return object reference. */
|
||||
sscanf(cursor, "ret:%d", &index);
|
||||
target = returns;
|
||||
}
|
||||
cursor = strtok(NULL, " ");
|
||||
CHECK(cursor);
|
||||
if (index >= utarray_len(target)) {
|
||||
utarray_resize(target, index + 1);
|
||||
}
|
||||
object_id *id = (object_id *) utarray_eltptr(target, index);
|
||||
hex_to_sha1(cursor, &id->id[0]);
|
||||
cursor = strtok(NULL, " ");
|
||||
}
|
||||
/* TODO(pcm): Implement pass by value. */
|
||||
/* Now assemble the task specification. */
|
||||
task_spec *spec =
|
||||
alloc_task_spec(function_id, utarray_len(args), utarray_len(returns), 0);
|
||||
for (int i = 0; i < utarray_len(args); ++i) {
|
||||
object_id *id = (object_id *) utarray_eltptr(args, i);
|
||||
task_args_add_ref(spec, *id);
|
||||
}
|
||||
for (int i = 0; i < utarray_len(returns); ++i) {
|
||||
object_id *id = (object_id *) utarray_eltptr(returns, i);
|
||||
*task_return(spec, i) = *id;
|
||||
}
|
||||
utarray_free(args);
|
||||
utarray_free(returns);
|
||||
return spec;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user