From 6a73711888196f6bce3e8d86a33e24545c41924b Mon Sep 17 00:00:00 2001 From: Stephanie Wang Date: Tue, 20 Dec 2016 00:13:39 -0800 Subject: [PATCH] Update the task table (#129) * Update the task table * Move updating task table out of scheduling algorithm. --- src/photon/photon.h | 2 ++ src/photon/photon_algorithm.c | 10 ++++++++++ src/photon/photon_scheduler.c | 35 ++++++++++++++++++++++------------- 3 files changed, 34 insertions(+), 13 deletions(-) diff --git a/src/photon/photon.h b/src/photon/photon.h index 29418004f..42b2b1b3e 100644 --- a/src/photon/photon.h +++ b/src/photon/photon.h @@ -31,6 +31,8 @@ enum photon_message_type { /** Contains all information that is associated to a worker. */ typedef struct { int sock; + /** A pointer to a task object, to update the task table. */ + task *task_in_progress; } worker; // clang-format on diff --git a/src/photon/photon_algorithm.c b/src/photon/photon_algorithm.c index acebca71f..22ad92abd 100644 --- a/src/photon/photon_algorithm.c +++ b/src/photon/photon_algorithm.c @@ -215,6 +215,12 @@ void queue_task_locally(local_scheduler_state *state, memcpy(elt->spec, spec, task_spec_size(spec)); elt->from_global_scheduler = from_global_scheduler; DL_APPEND(algorithm_state->task_queue, elt); + if (!from_global_scheduler && state->db != NULL) { + task *task = + alloc_task(spec, TASK_STATUS_SCHEDULED, get_db_client_id(state->db)); + task_table_add_task(state->db, task, (retry_info *) &photon_retry, NULL, + NULL); + } } void give_task_to_global_scheduler(local_scheduler_state *state, @@ -282,6 +288,10 @@ void handle_task_scheduled(local_scheduler_state *state, void handle_worker_available(local_scheduler_state *state, scheduling_algorithm_state *algorithm_state, int worker_index) { + worker *available_worker = + (worker *) utarray_eltptr(state->workers, worker_index); + CHECK(available_worker->task_in_progress == NULL); + /* Try to schedule another task to the worker. */ int scheduled_task = find_and_schedule_task_if_possible(state, algorithm_state, worker_index); /* If we couldn't find a task to schedule, add the worker to the queue of diff --git a/src/photon/photon_scheduler.c b/src/photon/photon_scheduler.c index 6b31f77a8..2647daf0e 100644 --- a/src/photon/photon_scheduler.c +++ b/src/photon/photon_scheduler.c @@ -89,18 +89,14 @@ void assign_task_to_worker(local_scheduler_state *state, write_message(w->sock, EXECUTE_TASK, task_spec_size(spec), (uint8_t *) spec); /* Update the global task table. */ if (state->db != NULL) { - retry_info retry; - memset(&retry, 0, sizeof(retry)); - retry.num_retries = 0; - retry.timeout = 100; - retry.fail_callback = NULL; task *task = alloc_task(spec, TASK_STATUS_RUNNING, get_db_client_id(state->db)); - if (from_global_scheduler) { - task_table_update(state->db, task, (retry_info *) &retry, NULL, NULL); - } else { - task_table_add_task(state->db, task, (retry_info *) &retry, NULL, NULL); - } + task_table_update(state->db, task, (retry_info *) &photon_retry, NULL, + NULL); + /* Record which task this worker is executing. This will be freed in + * process_message when the worker sends a GET_TASK message to the local + * scheduler. */ + w->task_in_progress = copy_task(task); } } @@ -142,11 +138,11 @@ void reconstruct_object_task_lookup_callback(object_id reconstruct_object_id, * reconstruction operation. NOTE: This codepath is not responsible for * detecting failure of the other reconstruction, or updating the * scheduling_state accordingly. */ - /* TODO(swang): Once we add code to modify the task table properly, this - * should also include TASK_STATUS_RUNNING. */ scheduling_state task_status = task_state(task); if (task_status == TASK_STATUS_WAITING || - task_status == TASK_STATUS_SCHEDULED) { + task_status == TASK_STATUS_SCHEDULED || + task_status == TASK_STATUS_RUNNING) { + LOG_DEBUG("Task to reconstruct had scheduling state %d", task_status); return; } /* Recursively reconstruct the task's inputs, if necessary. */ @@ -207,6 +203,19 @@ void process_message(event_loop *loop, case GET_TASK: { worker_index *wi; HASH_FIND_INT(state->worker_index, &client_sock, wi); + /* Update the task table with the completed task. */ + worker *available_worker = + (worker *) utarray_eltptr(state->workers, wi->worker_index); + if (state->db != NULL && available_worker->task_in_progress != NULL) { + task_set_state(available_worker->task_in_progress, TASK_STATUS_DONE); + task_table_update(state->db, available_worker->task_in_progress, + (retry_info *) &photon_retry, NULL, NULL); + /* The call to task_table_update takes ownership of the task_in_progress, + * so we set the pointer to NULL so it is not used. */ + available_worker->task_in_progress = NULL; + } + /* Let the scheduling algorithm process the fact that there is an available + * worker. */ handle_worker_available(state, state->algorithm_state, wi->worker_index); } break; case RECONSTRUCT_OBJECT: {