diff --git a/src/ray/raylet/scheduling/cluster_task_manager.cc b/src/ray/raylet/scheduling/cluster_task_manager.cc index af34d1dd5..baa467ae8 100644 --- a/src/ray/raylet/scheduling/cluster_task_manager.cc +++ b/src/ray/raylet/scheduling/cluster_task_manager.cc @@ -31,6 +31,8 @@ bool ClusterTaskManager::SchedulePendingTasks() { auto request_resources = task.GetTaskSpecification().GetRequiredResources().GetResourceMap(); int64_t _unused; + // TODO (Alex): We should distinguish between infeasible tasks and a fully + // utilized cluster. std::string node_id_string = cluster_resource_scheduler_->GetBestSchedulableNode(request_resources, &_unused); if (node_id_string.empty()) { diff --git a/src/ray/raylet/scheduling/cluster_task_manager_test.cc b/src/ray/raylet/scheduling/cluster_task_manager_test.cc index 290cd7fa3..6133c0dae 100644 --- a/src/ray/raylet/scheduling/cluster_task_manager_test.cc +++ b/src/ray/raylet/scheduling/cluster_task_manager_test.cc @@ -250,7 +250,8 @@ std::shared_ptr CreateSingleNodeScheduler( return scheduler; } -Task CreateTask(const std::unordered_map &required_resources) { +Task CreateTask(const std::unordered_map &required_resources, + int num_args = 0) { TaskSpecBuilder spec_builder; TaskID id = RandomTaskId(); JobID job_id = RandomJobId(); @@ -258,6 +259,12 @@ Task CreateTask(const std::unordered_map &required_resource spec_builder.SetCommonTaskSpec( id, Language::PYTHON, FunctionDescriptorBuilder::BuildPython("", "", "", ""), job_id, TaskID::Nil(), 0, TaskID::Nil(), address, 0, required_resources, {}); + + for (int i = 0; i < num_args; i++) { + ObjectID put_id = ObjectID::ForPut(TaskID::Nil(), /*index=*/i + 1); + spec_builder.AddArg(TaskArgByReference(put_id, rpc::Address())); + } + rpc::TaskExecutionSpec execution_spec_message; execution_spec_message.set_num_forwards(1); return Task(spec_builder.Build(), TaskExecutionSpecification(execution_spec_message)); @@ -326,13 +333,102 @@ TEST_F(ClusterTaskManagerTest, BasicTest) { task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_); - ASSERT_EQ(callback_occurred, true); + ASSERT_TRUE(callback_occurred); ASSERT_EQ(leased_workers_.size(), 1); ASSERT_EQ(pool_.workers.size(), 0); ASSERT_EQ(fulfills_dependencies_calls_, 0); ASSERT_EQ(node_info_calls_, 0); } +TEST_F(ClusterTaskManagerTest, NoFeasibleNodeTest) { + std::shared_ptr worker = + std::make_shared(WorkerID::FromRandom(), 1234); + pool_.PushWorker(std::dynamic_pointer_cast(worker)); + + Task task = CreateTask({{ray::kCPU_ResourceLabel, 999}}); + rpc::RequestWorkerLeaseReply reply; + + bool callback_called = false; + bool *callback_called_ptr = &callback_called; + auto callback = [callback_called_ptr]() { *callback_called_ptr = true; }; + + task_manager_.QueueTask(task, &reply, callback); + task_manager_.SchedulePendingTasks(); + task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_); + + ASSERT_FALSE(callback_called); + ASSERT_EQ(leased_workers_.size(), 0); + // Worker is unused. + ASSERT_EQ(pool_.workers.size(), 1); + ASSERT_EQ(fulfills_dependencies_calls_, 0); + ASSERT_EQ(node_info_calls_, 0); +} + +TEST_F(ClusterTaskManagerTest, ResourceTakenWhileResolving) { + /* + Test the race condition in which a task is assigned to a node, but cannot + run because its dependencies are unresolved. Once its dependencies are + resolved, the node no longer has available resources. + */ + std::shared_ptr worker = + std::make_shared(WorkerID::FromRandom(), 1234); + std::shared_ptr worker2 = + std::make_shared(WorkerID::FromRandom(), 12345); + pool_.PushWorker(std::dynamic_pointer_cast(worker2)); + pool_.PushWorker(std::dynamic_pointer_cast(worker)); + + rpc::RequestWorkerLeaseReply reply; + int num_callbacks = 0; + int *num_callbacks_ptr = &num_callbacks; + auto callback = [num_callbacks_ptr]() { + (*num_callbacks_ptr) = *num_callbacks_ptr + 1; + }; + + /* Blocked on dependencies */ + auto task = CreateTask({{ray::kCPU_ResourceLabel, 5}}, 1); + dependencies_fulfilled_ = false; + task_manager_.QueueTask(task, &reply, callback); + task_manager_.SchedulePendingTasks(); + task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_); + + ASSERT_EQ(num_callbacks, 0); + ASSERT_EQ(leased_workers_.size(), 0); + ASSERT_EQ(pool_.workers.size(), 2); + + /* This task can run */ + auto task2 = CreateTask({{ray::kCPU_ResourceLabel, 5}}); + task_manager_.QueueTask(task2, &reply, callback); + task_manager_.SchedulePendingTasks(); + task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_); + + ASSERT_EQ(num_callbacks, 1); + ASSERT_EQ(leased_workers_.size(), 1); + ASSERT_EQ(pool_.workers.size(), 1); + + /* First task is unblocked now, but resources are no longer available */ + auto id = task.GetTaskSpecification().TaskId(); + std::vector unblocked = {id}; + task_manager_.TasksUnblocked(unblocked); + task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_); + + ASSERT_EQ(num_callbacks, 1); + ASSERT_EQ(leased_workers_.size(), 1); + ASSERT_EQ(pool_.workers.size(), 1); + + /* Second task finishes, making space for the original task */ + single_node_resource_scheduler_->FreeLocalTaskResources( + worker->GetAllocatedInstances()); + // single_node_resource_scheduler_->UpdateLocalAvailableResourcesFromResourceInstances(); + leased_workers_.clear(); + + task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_); + + // Task2 is now done so task can run. + ASSERT_EQ(num_callbacks, 2); + ASSERT_EQ(leased_workers_.size(), 1); + ASSERT_EQ(pool_.workers.size(), 0); +} + } // namespace raylet } // namespace ray