#include #include #include #include #include #include "daggy/loggers/dag_run/OStreamLogger.hpp" using namespace daggy; using namespace daggy::loggers::dag_run; const TaskSet SAMPLE_TASKS{ {"work_a", Task{.job{{"command", std::vector{"/bin/echo", "a"}}}, .children{"c"}}}, {"work_b", Task{.job{{"command", std::vector{"/bin/echo", "b"}}}, .children{"c"}}}, {"work_c", Task{.job{{"command", std::vector{"/bin/echo", "c"}}}}}}; inline DAGRunID testDAGRunInit(DAGRunLogger &logger, const std::string &tag, const TaskSet &tasks) { auto runID = logger.startDAGRun(DAGSpec{.tag = tag, .tasks = tasks}); // Verify run shows up in the list { auto runs = logger.queryDAGRuns(); REQUIRE(!runs.empty()); auto it = std::find_if(runs.begin(), runs.end(), [runID](const auto &r) { return r.runID == runID; }); REQUIRE(it != runs.end()); REQUIRE(it->tag == tag); REQUIRE(it->runState == +RunState::QUEUED); } // Verify states { REQUIRE(logger.getDAGRunState(runID) == +RunState::QUEUED); for (const auto &[k, _] : tasks) { REQUIRE(logger.getTaskState(runID, k) == +RunState::QUEUED); } } // Verify integrity of run { auto dagRun = logger.getDAGRun(runID); REQUIRE(dagRun.dagSpec.tag == tag); REQUIRE(dagRun.dagSpec.tasks == tasks); REQUIRE(dagRun.taskRunStates.size() == tasks.size()); auto nonQueuedTask = std::find_if( dagRun.taskRunStates.begin(), dagRun.taskRunStates.end(), [](const auto &a) { return a.second != +RunState::QUEUED; }); REQUIRE(nonQueuedTask == dagRun.taskRunStates.end()); REQUIRE(dagRun.dagStateChanges.size() == 1); REQUIRE(dagRun.dagStateChanges.back().newState == +RunState::QUEUED); } // Update DAG state and ensure that it's updated; { logger.updateDAGRunState(runID, RunState::RUNNING); auto dagRun = logger.getDAGRun(runID); REQUIRE(dagRun.dagStateChanges.back().newState == +RunState::RUNNING); } // Update a task state { for (const auto &[k, v] : tasks) logger.updateTaskState(runID, k, RunState::RUNNING); auto dagRun = logger.getDAGRun(runID); for (const auto &[k, v] : tasks) { REQUIRE(dagRun.taskRunStates.at(k) == +RunState::RUNNING); } } return runID; } TEST_CASE("ostream_logger", "[ostream_logger]") { std::stringstream ss; daggy::loggers::dag_run::OStreamLogger logger(ss); SECTION("DAGRun Starts") { testDAGRunInit(logger, "init_test", SAMPLE_TASKS); } }