[ Web Proxy ]
URL:
Viewing: https://raw.githubusercontent.com/ModuleWorks/taskflow/task_isolation/unittests/test_pipelines.cpp [Back]  [Original]

#define DOCTEST_CONFIG_IMPLEMENT_WITH_MAIN

#include 

#include 
#include 

#include      /* srand, rand */
#include        /* time */

// --------------------------------------------------------
// Testcase: 1 pipe, L lines, w workers
// --------------------------------------------------------
void pipeline_1P(size_t L, unsigned w, tf::PipeType type) {

  tf::Executor executor(w);

  const size_t maxN = 100;

  std::vector source(maxN);
  std::iota(source.begin(), source.end(), 0);

  // iterate different data amount (1, 2, 3, 4, 5, ... 1000000)
  for (size_t N = 0; N = N) {
    //      pf.stop();
    //      return;
    //    }
    //    std::scoped_lock lock(mutex);
    //    collection.push_back(ticket);
    //  }});

    //  taskflow.composed_of(pl);
    //  executor.run(taskflow).wait();
    //  REQUIRE(collection.size() == N);
    //  std::sort(collection.begin(), collection.end());
    //  for(size_t k=0; k O --
//
// ----------------------------------------------------------------------------

void three_parallel_pipelines(size_t L, unsigned w) {

  tf::Executor executor(w);

  const size_t maxN = 100;

  std::vector source(maxN);
  std::iota(source.begin(), source.end(), 0);
  std::vector mybuffer1(L);
  std::vector mybuffer2(L);
  std::vector mybuffer3(L);

  for(size_t N = 0; N  SSSS -> O -> SSP -> O -> SP -> O
//
// ----------------------------------------------------------------------------

void three_concatenated_pipelines(size_t L, unsigned w) {

  tf::Executor executor(w);

  const size_t maxN = 100;

  std::vector source(maxN);
  std::iota(source.begin(), source.end(), 0);
  std::vector mybuffer1(L);
  std::vector mybuffer2(L);
  std::vector mybuffer3(L);

  for(size_t N = 0; N  SPSP -> conditional_task
//        ^            |
//        |____________|
// ----------------------------------------------------------------------------

void looping_pipelines(size_t L, unsigned w) {

  tf::Executor executor(w);

  const size_t maxN = 100;

  std::vector source(maxN);
  std::iota(source.begin(), source.end(), 0);
  std::vector mybuffer(L);

  tf::Taskflow taskflow;

  size_t j1 = 0, j3 = 0;
  std::atomic j2 = 0;
  std::atomic j4 = 0;
  std::mutex mutex2;
  std::mutex mutex4;
  std::vector collection2;
  std::vector collection4;
  size_t cnt = 0;

  size_t N = 0;

  tf::Pipeline pl(L,
    tf::Pipe{tf::PipeType::SERIAL, [&N, &source, &j1, &mybuffer, L](auto& pf) mutable {
      if(j1 == N) {
        pf.stop();
        return;
      }
      REQUIRE(j1 == source[j1]);
      REQUIRE(pf.token() % L == pf.line());
      mybuffer[pf.line()][pf.pipe()] = source[j1] + 1;
      j1++;
    }},

    tf::Pipe{tf::PipeType::PARALLEL, [&N, &j2, &mutex2, &collection2, &mybuffer, L](auto& pf) mutable {
      REQUIRE(j2++ < N);
      {
        std::scoped_lock lock(mutex2);
        REQUIRE(pf.token() % L == pf.line());
        mybuffer[pf.line()][pf.pipe()] = mybuffer[pf.line()][pf.pipe() - 1] + 1;
        collection2.push_back(mybuffer[pf.line()][pf.pipe() - 1]);
      }
    }},

    tf::Pipe{tf::PipeType::SERIAL, [&N, &source, &j3, &mybuffer, L](auto& pf) mutable {
      REQUIRE(j3 < N);
      REQUIRE(pf.token() % L == pf.line());
      REQUIRE(source[j3] + 2 == mybuffer[pf.line()][pf.pipe() - 1]);
      mybuffer[pf.line()][pf.pipe()] = mybuffer[pf.line()][pf.pipe() - 1] + 1;
      j3++;
    }},

    tf::Pipe{tf::PipeType::PARALLEL, [&N, &j4, &mutex4, &collection4, &mybuffer, L](auto& pf) mutable {
      REQUIRE(j4++ < N);
      {
        std::scoped_lock lock(mutex4);
        REQUIRE(pf.token() % L == pf.line());
        collection4.push_back(mybuffer[pf.line()][pf.pipe() - 1]);
      }
    }}
  );

  auto pipeline = taskflow.composed_of(pl).name("module_of_pipeline");
  auto initial = taskflow.emplace([](){}).name("initial");

  auto conditional = taskflow.emplace([&](){
    REQUIRE(j1 == N);
    REQUIRE(j2 == N);
    REQUIRE(j3 == N);
    REQUIRE(j4 == N);
    REQUIRE(collection2.size() == N);
    REQUIRE(collection4.size() == N);
    std::sort(collection2.begin(), collection2.end());
    std::sort(collection4.begin(), collection4.end());
    for (size_t i = 0; i < N; ++i) {
      REQUIRE(collection2[i] == i + 1);
      REQUIRE(collection4[i] == i + 3);
    }
    REQUIRE(pl.num_tokens() == cnt);

    // reset variables
    j1 = j2 = j3 = j4 = 0;
    for(size_t i = 0; i < mybuffer.size(); ++i){
      for(size_t j = 0; j < mybuffer[0].size(); ++j){
        mybuffer[i][j] = 0;
      }
    }
    collection2.clear();
    collection4.clear();
    ++N;
    cnt+=N;

    return N < maxN ? 0 : 1;
  }).name("conditional");

  auto terminal = taskflow.emplace([](){}).name("terminal");

  initial.precede(pipeline);
  pipeline.precede(conditional);
  conditional.precede(pipeline, terminal);

  executor.run(taskflow).wait();
}

// looping piplines
TEST_CASE("Looping.Pipelines.1L.1W" * doctest::timeout(300)) {
  looping_pipelines(1, 1);
}

TEST_CASE("Looping.Pipelines.1L.2W" * doctest::timeout(300)) {
  looping_pipelines(1, 2);
}

TEST_CASE("Looping.Pipelines.1L.3W" * doctest::timeout(300)) {
  looping_pipelines(1, 3);
}

TEST_CASE("Looping.Pipelines.1L.4W" * doctest::timeout(300)) {
  looping_pipelines(1, 4);
}

TEST_CASE("Looping.Pipelines.1L.5W" * doctest::timeout(300)) {
  looping_pipelines(1, 5);
}

TEST_CASE("Looping.Pipelines.1L.6W" * doctest::timeout(300)) {
  looping_pipelines(1, 6);
}

TEST_CASE("Looping.Pipelines.1L.7W" * doctest::timeout(300)) {
  looping_pipelines(1, 7);
}

TEST_CASE("Looping.Pipelines.1L.8W" * doctest::timeout(300)) {
  looping_pipelines(1, 8);
}

TEST_CASE("Looping.Pipelines.2L.1W" * doctest::timeout(300)) {
  looping_pipelines(2, 1);
}

TEST_CASE("Looping.Pipelines.2L.2W" * doctest::timeout(300)) {
  looping_pipelines(2, 2);
}

TEST_CASE("Looping.Pipelines.2L.3W" * doctest::timeout(300)) {
  looping_pipelines(2, 3);
}

TEST_CASE("Looping.Pipelines.2L.4W" * doctest::timeout(300)) {
  looping_pipelines(2, 4);
}

TEST_CASE("Looping.Pipelines.2L.5W" * doctest::timeout(300)) {
  looping_pipelines(2, 5);
}

TEST_CASE("Looping.Pipelines.2L.6W" * doctest::timeout(300)) {
  looping_pipelines(2, 6);
}

TEST_CASE("Looping.Pipelines.2L.7W" * doctest::timeout(300)) {
  looping_pipelines(2, 7);
}

TEST_CASE("Looping.Pipelines.2L.8W" * doctest::timeout(300)) {
  looping_pipelines(2, 8);
}

TEST_CASE("Looping.Pipelines.3L.1W" * doctest::timeout(300)) {
  looping_pipelines(3, 1);
}

TEST_CASE("Looping.Pipelines.3L.2W" * doctest::timeout(300)) {
  looping_pipelines(3, 2);
}

TEST_CASE("Looping.Pipelines.3L.3W" * doctest::timeout(300)) {
  looping_pipelines(3, 3);
}

TEST_CASE("Looping.Pipelines.3L.4W" * doctest::timeout(300)) {
  looping_pipelines(3, 4);
}

TEST_CASE("Looping.Pipelines.3L.5W" * doctest::timeout(300)) {
  looping_pipelines(3, 5);
}

TEST_CASE("Looping.Pipelines.3L.6W" * doctest::timeout(300)) {
  looping_pipelines(3, 6);
}

TEST_CASE("Looping.Pipelines.3L.7W" * doctest::timeout(300)) {
  looping_pipelines(3, 7);
}

TEST_CASE("Looping.Pipelines.3L.8W" * doctest::timeout(300)) {
  looping_pipelines(3, 8);
}

TEST_CASE("Looping.Pipelines.4L.1W" * doctest::timeout(300)) {
  looping_pipelines(4, 1);
}

TEST_CASE("Looping.Pipelines.4L.2W" * doctest::timeout(300)) {
  looping_pipelines(4, 2);
}

TEST_CASE("Looping.Pipelines.4L.3W" * doctest::timeout(300)) {
  looping_pipelines(4, 3);
}

TEST_CASE("Looping.Pipelines.4L.4W" * doctest::timeout(300)) {
  looping_pipelines(4, 4);
}

TEST_CASE("Looping.Pipelines.4L.5W" * doctest::timeout(300)) {
  looping_pipelines(4, 5);
}

TEST_CASE("Looping.Pipelines.4L.6W" * doctest::timeout(300)) {
  looping_pipelines(4, 6);
}

TEST_CASE("Looping.Pipelines.4L.7W" * doctest::timeout(300)) {
  looping_pipelines(4, 7);
}

TEST_CASE("Looping.Pipelines.4L.8W" * doctest::timeout(300)) {
  looping_pipelines(4, 8);
}

TEST_CASE("Looping.Pipelines.5L.1W" * doctest::timeout(300)) {
  looping_pipelines(5, 1);
}

TEST_CASE("Looping.Pipelines.5L.2W" * doctest::timeout(300)) {
  looping_pipelines(5, 2);
}

TEST_CASE("Looping.Pipelines.5L.3W" * doctest::timeout(300)) {
  looping_pipelines(5, 3);
}

TEST_CASE("Looping.Pipelines.5L.4W" * doctest::timeout(300)) {
  looping_pipelines(5, 4);
}

TEST_CASE("Looping.Pipelines.5L.5W" * doctest::timeout(300)) {
  looping_pipelines(5, 5);
}

TEST_CASE("Looping.Pipelines.5L.6W" * doctest::timeout(300)) {
  looping_pipelines(5, 6);
}

TEST_CASE("Looping.Pipelines.5L.7W" * doctest::timeout(300)) {
  looping_pipelines(5, 7);
}

TEST_CASE("Looping.Pipelines.5L.8W" * doctest::timeout(300)) {
  looping_pipelines(5, 8);
}

TEST_CASE("Looping.Pipelines.6L.1W" * doctest::timeout(300)) {
  looping_pipelines(6, 1);
}

TEST_CASE("Looping.Pipelines.6L.2W" * doctest::timeout(300)) {
  looping_pipelines(6, 2);
}

TEST_CASE("Looping.Pipelines.6L.3W" * doctest::timeout(300)) {
  looping_pipelines(6, 3);
}

TEST_CASE("Looping.Pipelines.6L.4W" * doctest::timeout(300)) {
  looping_pipelines(6, 4);
}

TEST_CASE("Looping.Pipelines.6L.5W" * doctest::timeout(300)) {
  looping_pipelines(6, 5);
}

TEST_CASE("Looping.Pipelines.6L.6W" * doctest::timeout(300)) {
  looping_pipelines(6, 6);
}

TEST_CASE("Looping.Pipelines.6L.7W" * doctest::timeout(300)) {
  looping_pipelines(6, 7);
}

TEST_CASE("Looping.Pipelines.6L.8W" * doctest::timeout(300)) {
  looping_pipelines(6, 8);
}

TEST_CASE("Looping.Pipelines.7L.1W" * doctest::timeout(300)) {
  looping_pipelines(7, 1);
}

TEST_CASE("Looping.Pipelines.7L.2W" * doctest::timeout(300)) {
  looping_pipelines(7, 2);
}

TEST_CASE("Looping.Pipelines.7L.3W" * doctest::timeout(300)) {
  looping_pipelines(7, 3);
}

TEST_CASE("Looping.Pipelines.7L.4W" * doctest::timeout(300)) {
  looping_pipelines(7, 4);
}

TEST_CASE("Looping.Pipelines.7L.5W" * doctest::timeout(300)) {
  looping_pipelines(7, 5);
}

TEST_CASE("Looping.Pipelines.7L.6W" * doctest::timeout(300)) {
  looping_pipelines(7, 6);
}

TEST_CASE("Looping.Pipelines.7L.7W" * doctest::timeout(300)) {
  looping_pipelines(7, 7);
}

TEST_CASE("Looping.Pipelines.7L.8W" * doctest::timeout(300)) {
  looping_pipelines(7, 8);
}

TEST_CASE("Looping.Pipelines.8L.1W" * doctest::timeout(300)) {
  looping_pipelines(8, 1);
}

TEST_CASE("Looping.Pipelines.8L.2W" * doctest::timeout(300)) {
  looping_pipelines(8, 2);
}

TEST_CASE("Looping.Pipelines.8L.3W" * doctest::timeout(300)) {
  looping_pipelines(8, 3);
}

TEST_CASE("Looping.Pipelines.8L.4W" * doctest::timeout(300)) {
  looping_pipelines(8, 4);
}

TEST_CASE("Looping.Pipelines.8L.5W" * doctest::timeout(300)) {
  looping_pipelines(8, 5);
}

TEST_CASE("Looping.Pipelines.8L.6W" * doctest::timeout(300)) {
  looping_pipelines(8, 6);
}

TEST_CASE("Looping.Pipelines.8L.7W" * doctest::timeout(300)) {
  looping_pipelines(8, 7);
}

TEST_CASE("Looping.Pipelines.8L.8W" * doctest::timeout(300)) {
  looping_pipelines(8, 8);
}

// ----------------------------------------------------------------------------
//
// ifelse pipeline has three pipes, L lines, w workers
//
// SPS
// ----------------------------------------------------------------------------

int ifelse_pipe_ans(int a) {
  // pipe 1
  if(a / 2 != 0) {
    a += 8;
  }
  // pipe 2
  if(a > 4897) {
    a -= 1834;
  }
  else {
    a += 3;
  }
  // pipe 3
  if((a + 9) / 4 < 50) {
    a += 1;
  }
  else {
    a += 17;
  }

  return a;
}

void ifelse_pipeline(size_t L, unsigned w) {
  //srand(time(NULL));

  tf::Executor executor(w);
  size_t maxN = 200;

  std::vector source(maxN);
  for(auto&& s: source) {
    s = rand() % 9962;
  }
  std::vector buffer(L);

  for(size_t N = 1; N < maxN; ++N) {
    tf::Taskflow taskflow;

    std::vector collection;
    collection.reserve(N);

    tf::Pipeline pl(L,
      // pipe 1
      tf::Pipe(tf::PipeType::SERIAL, [&, N](auto& pf){
        if(pf.token() == N) {
          pf.stop();
          return;
        }

        if(source[pf.token()] / 2 == 0) {
          buffer[pf.line()][pf.pipe()] = source[pf.token()];
        }
        else {
          buffer[pf.line()][pf.pipe()] = source[pf.token()] + 8;
        }

      }),

      // pipe 2
      tf::Pipe(tf::PipeType::PARALLEL, [&](auto& pf){

        if(buffer[pf.line()][pf.pipe() - 1] > 4897) {
          buffer[pf.line()][pf.pipe()] =  buffer[pf.line()][pf.pipe() - 1] - 1834;
        }
        else {
          buffer[pf.line()][pf.pipe()] = buffer[pf.line()][pf.pipe() - 1] + 3;
        }

      }),

      // pipe 3
      tf::Pipe(tf::PipeType::SERIAL, [&](auto& pf){

        if((buffer[pf.line()][pf.pipe() - 1] + 9) / 4 < 50) {
          buffer[pf.line()][pf.pipe()] = buffer[pf.line()][pf.pipe() - 1] + 1;
        }
        else {
          buffer[pf.line()][pf.pipe()] = buffer[pf.line()][pf.pipe() - 1] + 17;
        }

        collection.push_back(buffer[pf.line()][pf.pipe()]);

      })
    );
    auto pl_t = taskflow.composed_of(pl).name("pipeline");

    auto check_t = taskflow.emplace([&](){
      for(size_t n = 0; n < N; ++n) {
        REQUIRE(collection[n] == ifelse_pipe_ans(source[n]));
      }
    }).name("check");

    pl_t.precede(check_t);

    executor.run(taskflow).wait();

  }
}

TEST_CASE("Ifelse.Pipelines.1L.1W" * doctest::timeout(300)) {
  ifelse_pipeline(1, 1);
}

TEST_CASE("Ifelse.Pipelines.1L.2W" * doctest::timeout(300)) {
  ifelse_pipeline(1, 2);
}

TEST_CASE("Ifelse.Pipelines.1L.3W" * doctest::timeout(300)) {
  ifelse_pipeline(1, 3);
}

TEST_CASE("Ifelse.Pipelines.1L.4W" * doctest::timeout(300)) {
  ifelse_pipeline(1, 4);
}

TEST_CASE("Ifelse.Pipelines.3L.1W" * doctest::timeout(300)) {
  ifelse_pipeline(3, 1);
}

TEST_CASE("Ifelse.Pipelines.3L.2W" * doctest::timeout(300)) {
  ifelse_pipeline(3, 2);
}

TEST_CASE("Ifelse.Pipelines.3L.3W" * doctest::timeout(300)) {
  ifelse_pipeline(3, 3);
}

TEST_CASE("Ifelse.Pipelines.3L.4W" * doctest::timeout(300)) {
  ifelse_pipeline(3, 4);
}

TEST_CASE("Ifelse.Pipelines.5L.1W" * doctest::timeout(300)) {
  ifelse_pipeline(5, 1);
}

TEST_CASE("Ifelse.Pipelines.5L.2W" * doctest::timeout(300)) {
  ifelse_pipeline(5, 2);
}

TEST_CASE("Ifelse.Pipelines.5L.3W" * doctest::timeout(300)) {
  ifelse_pipeline(5, 3);
}

TEST_CASE("Ifelse.Pipelines.5L.4W" * doctest::timeout(300)) {
  ifelse_pipeline(5, 4);
}

TEST_CASE("Ifelse.Pipelines.7L.1W" * doctest::timeout(300)) {
  ifelse_pipeline(7, 1);
}

TEST_CASE("Ifelse.Pipelines.7L.2W" * doctest::timeout(300)) {
  ifelse_pipeline(7, 2);
}

TEST_CASE("Ifelse.Pipelines.7L.3W" * doctest::timeout(300)) {
  ifelse_pipeline(7, 3);
}

TEST_CASE("Ifelse.Pipelines.7L.4W" * doctest::timeout(300)) {
  ifelse_pipeline(7, 4);
}

// ----------------------------------------------------------------------------
// pipeline in pipeline
// pipeline has 4 pipes, L lines, W workers
// each subpipeline has 3 pipes, subL lines
//
// pipeline = SPPS
// each subpipeline = SPS
//
// ----------------------------------------------------------------------------

void pipeline_in_pipeline(size_t L, unsigned w, unsigned subL) {


  tf::Executor executor(w);

  const size_t maxN = 5;
  const size_t maxsubN = 4;

  std::vector source(maxN);
  for(auto&& each: source) {
    each.resize(maxsubN);
    std::iota(each.begin(), each.end(), 0);
  }

  std::vector buffer(L);

  // each pipe contains one subpipeline
  // each subpipeline has three pipes, subL lines
  //
  // subbuffers[0][1][2][2] means
  // first line, second pipe, third subline, third subpipe
  std::vector subbuffers(L);

  for(auto&& pipes: subbuffers) {
    pipes.resize(4);
    for(auto&& pipe: pipes) {
        pipe.resize(subL);
    }
  }

  for (size_t N = 1; N < maxN; ++N) {
    for(size_t subN = 1; subN < maxsubN; ++subN) {

      size_t j1 = 0, j4 = 0;
      std::atomic j2 = 0;
      std::atomic j3 = 0;

      // begin of pipeline ---------------------------
      tf::Pipeline pl(L,

        // begin of pipe 1 -----------------------------
        tf::Pipe{tf::PipeType::SERIAL, [&, N, subN, subL](auto& pf) mutable {
          if(j1 == N) {
            pf.stop();
            return;
          }

          size_t subj1 = 0, subj3 = 0;
          std::atomic subj2 = 0;
          std::vector subcollection;
          subcollection.reserve(subN);

          // subpipeline
          tf::Pipeline subpl(subL,

            // subpipe 1
            tf::Pipe{tf::PipeType::SERIAL, [&, subN](auto& subpf) mutable {
              if(subj1 == subN) {
                subpf.stop();
                return;
              }

              REQUIRE(subpf.token() % subL == subpf.line());

              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subj1] + 1;

              ++subj1;
            }},

            // subpipe 2
            tf::Pipe{tf::PipeType::PARALLEL, [&, subN](auto& subpf) mutable {
              REQUIRE(subj2++ < subN);
              REQUIRE(subpf.token() % subL == subpf.line());
              REQUIRE(source[pf.token()][subpf.token()] + 1 == subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe() - 1]);
              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subpf.token()] + 1;
            }},


            // subpipe 3
            tf::Pipe{tf::PipeType::SERIAL, [&, subN](auto& subpf) mutable {
              REQUIRE(subj3 < subN);
              REQUIRE(subpf.token() % subL == subpf.line());
              REQUIRE(source[pf.token()][subj3] + 1 == subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe() - 1]);
              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subj3] + 3;
              subcollection.push_back(subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]);
              ++subj3;
            }}
          );

          tf::Taskflow taskflow;

          // test task
          auto test_t = taskflow.emplace([&, subN](){
            REQUIRE(subj1 == subN);
            REQUIRE(subj2 == subN);
            REQUIRE(subj3 == subN);
            //REQUIRE(subpl.num_tokens() == subN);
            REQUIRE(subcollection.size() == subN);
          }).name("test");

          // subpipeline
          auto subpl_t = taskflow.composed_of(subpl).name("module_of_subpipeline");

          subpl_t.precede(test_t);
          executor.corun(taskflow);

          buffer[pf.line()][pf.pipe()] = std::accumulate(
            subcollection.begin(),
            subcollection.end(),
            0
          );

          j1++;
        }},
        // end of pipe 1 -----------------------------

         //begin of pipe 2 ---------------------------
        tf::Pipe{tf::PipeType::PARALLEL, [&, N, subN, subL](auto& pf) mutable {

          REQUIRE(j2++ < N);
          int res = std::accumulate(
            source[pf.token()].begin(),
            source[pf.token()].begin() + subN,
            0
          );
          REQUIRE(buffer[pf.line()][pf.pipe() - 1] == res + 3 * subN);

          size_t subj1 = 0, subj3 = 0;
          std::atomic subj2 = 0;
          std::vector subcollection;
          subcollection.reserve(subN);

          // subpipeline
          tf::Pipeline subpl(subL,

            // subpipe 1
            tf::Pipe{tf::PipeType::SERIAL, [&, subN](auto& subpf) mutable {
              if(subj1 == subN) {
                subpf.stop();
                return;
              }

              REQUIRE(subpf.token() % subL == subpf.line());

              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subj1] + 1;

              ++subj1;
            }},

            // subpipe 2
            tf::Pipe{tf::PipeType::PARALLEL, [&, subN](auto& subpf) mutable {
              REQUIRE(subj2++ < subN);
              REQUIRE(subpf.token() % subL == subpf.line());
              REQUIRE(source[j2][subpf.token()] + 1 == subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe() - 1]);
              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subpf.token()] + 1;
            }},


            // subpipe 3
            tf::Pipe{tf::PipeType::SERIAL, [&, subN](auto& subpf) mutable {
              REQUIRE(subj3 < subN);
              REQUIRE(subpf.token() % subL == subpf.line());
              REQUIRE(source[pf.token()][subj3] + 1 == subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe() - 1]);
              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subj3] + 13;
              subcollection.push_back(subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]);
              ++subj3;
            }}
          );

          tf::Taskflow taskflow;

          // test task
          auto test_t = taskflow.emplace([&, subN](){
            REQUIRE(subj1 == subN);
            REQUIRE(subj2 == subN);
            REQUIRE(subj3 == subN);
            //REQUIRE(subpl.num_tokens() == subN);
            REQUIRE(subcollection.size() == subN);
          }).name("test");

          // subpipeline
          auto subpl_t = taskflow.composed_of(subpl).name("module_of_subpipeline");

          subpl_t.precede(test_t);
          executor.corun(taskflow);

          buffer[pf.line()][pf.pipe()] = std::accumulate(
            subcollection.begin(),
            subcollection.end(),
            0
          );

        }},
        // end of pipe 2 -----------------------------

        // begin of pipe 3 ---------------------------
        tf::Pipe{tf::PipeType::SERIAL, [&, N, subN, subL](auto& pf) mutable {

          REQUIRE(j3++ < N);
          int res = std::accumulate(
            source[pf.token()].begin(),
            source[pf.token()].begin() + subN,
            0
          );

          REQUIRE(buffer[pf.line()][pf.pipe() - 1] == res + 13 * subN);

          size_t subj1 = 0, subj3 = 0;
          std::atomic subj2 = 0;
          std::vector subcollection;
          subcollection.reserve(subN);

          // subpipeline
          tf::Pipeline subpl(subL,

            // subpipe 1
            tf::Pipe{tf::PipeType::SERIAL, [&, subN](auto& subpf) mutable {
              if(subj1 == subN) {
                subpf.stop();
                return;
              }

              REQUIRE(subpf.token() % subL == subpf.line());

              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subj1] + 1;

              ++subj1;
            }},

            // subpipe 2
            tf::Pipe{tf::PipeType::PARALLEL, [&, subN](auto& subpf) mutable {
              REQUIRE(subj2++ < subN);
              REQUIRE(subpf.token() % subL == subpf.line());
              REQUIRE(source[pf.token()][subpf.token()] + 1 == subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe() - 1]);
              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subpf.token()] + 1;
            }},


            // subpipe 3
            tf::Pipe{tf::PipeType::SERIAL, [&, subN](auto& subpf) mutable {
              REQUIRE(subj3 < subN);
              REQUIRE(subpf.token() % subL == subpf.line());
              REQUIRE(source[pf.token()][subj3] + 1 == subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe() - 1]);
              subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]
                = source[pf.token()][subj3] + 7;
              subcollection.push_back(subbuffers[pf.line()][pf.pipe()][subpf.line()][subpf.pipe()]);
              ++subj3;
            }}
          );

          tf::Taskflow taskflow;

          // test task
          auto test_t = taskflow.emplace([&, subN](){
            REQUIRE(subj1 == subN);
            REQUIRE(subj2 == subN);
            REQUIRE(subj3 == subN);
            //REQUIRE(subpl.num_tokens() == subN);
            REQUIRE(subcollection.size() == subN);
          }).name("test");

          // subpipeline
          auto subpl_t = taskflow.composed_of(subpl).name("module_of_subpipeline");

          subpl_t.precede(test_t);
          executor.corun(taskflow);

          buffer[pf.line()][pf.pipe()] = std::accumulate(
            subcollection.begin(),
            subcollection.end(),
            0
          );

        }},
        // end of pipe 3 -----------------------------

        // begin of pipe 4 ---------------------------
        tf::Pipe{tf::PipeType::SERIAL, [&, subN](auto& pf) mutable {

          int res = std::accumulate(
            source[j4].begin(),
            source[j4].begin() + subN,
            0
          );
          REQUIRE(buffer[pf.line()][pf.pipe() - 1] == res + 7 * subN);
          j4++;
        }}
        // end of pipe 4 -----------------------------
      );

      tf::Taskflow taskflow;
      taskflow.composed_of(pl).name("module_of_pipeline");
      executor.run(taskflow).wait();
    }
  }
}

TEST_CASE("PipelineinPipeline.Pipelines.1L.1W.1subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(1, 1, 1);
}

TEST_CASE("PipelineinPipeline.Pipelines.1L.1W.3subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(1, 1, 3);
}

TEST_CASE("PipelineinPipeline.Pipelines.1L.1W.4subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(1, 1, 4);
}

TEST_CASE("PipelineinPipeline.Pipelines.1L.2W.1subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(1, 2, 1);
}

TEST_CASE("PipelineinPipeline.Pipelines.1L.2W.3subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(1, 2, 3);
}

TEST_CASE("PipelineinPipeline.Pipelines.1L.2W.4subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(1, 2, 4);
}

TEST_CASE("PipelineinPipeline.Pipelines.3L.1W.1subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(3, 1, 1);
}

TEST_CASE("PipelineinPipeline.Pipelines.3L.1W.3subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(3, 1, 3);
}

TEST_CASE("PipelineinPipeline.Pipelines.3L.1W.4subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(3, 1, 4);
}

TEST_CASE("PipelineinPipeline.Pipelines.3L.2W.1subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(3, 2, 1);
}

TEST_CASE("PipelineinPipeline.Pipelines.3L.2W.3subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(3, 2, 3);
}

TEST_CASE("PipelineinPipeline.Pipelines.3L.2W.4subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(3, 2, 4);
}

TEST_CASE("PipelineinPipeline.Pipelines.5L.1W.1subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(5, 1, 1);
}

TEST_CASE("PipelineinPipeline.Pipelines.5L.1W.3subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(5, 1, 3);
}

TEST_CASE("PipelineinPipeline.Pipelines.5L.1W.4subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(5, 1, 4);
}

TEST_CASE("PipelineinPipeline.Pipelines.5L.2W.1subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(5, 2, 1);
}

TEST_CASE("PipelineinPipeline.Pipelines.5L.2W.3subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(5, 2, 3);
}

TEST_CASE("PipelineinPipeline.Pipelines.5L.2W.4subL" * doctest::timeout(300)) {
  pipeline_in_pipeline(5, 2, 4);
}

Web Proxy Viewer  |  New URL  |  Original Page