Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions bindings/python_bindings.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,8 @@ PYBIND11_MODULE(_py_capio_cl, m) {
});

m.def("serialize", &capiocl::serializer::Serializer::dump, py::arg("engine"),
py::arg("filename"), py::arg("version") = capiocl::CAPIO_CL_VERSION::V1);
py::arg("filename"), py::arg("compress") = false,
py::arg("version") = capiocl::CAPIO_CL_VERSION::V1);

py::class_<capiocl::engine::CapioCLEntry>(m, "CapioCLEntry")
.def(py::init<>())
Expand All @@ -150,4 +151,4 @@ PYBIND11_MODULE(_py_capio_cl, m) {
.def_readwrite("is_file", &capiocl::engine::CapioCLEntry::is_file)
.def_static("from_json", &capiocl::engine::CapioCLEntry::fromJson, py::arg("in"))
.def("to_json", &capiocl::engine::CapioCLEntry::toJson);
}
}
28 changes: 24 additions & 4 deletions capiocl/serializer.h
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
#ifndef CAPIO_CL_SERIALIZER_H
#define CAPIO_CL_SERIALIZER_H

#include <filesystem>
#include <utility>
#include <vector>

#include "capiocl.hpp"

/// @brief Namespace containing the CAPIO-CL Serializer component
Expand Down Expand Up @@ -28,6 +32,19 @@ class SerializerException final : public std::exception {
/// @brief Dump the current loaded CAPIO-CL configuration from class Engine to a CAPIO-CL
/// configuration file.
class Serializer final {
/**
* Compress entries from a CAPIO-CL engine into entries using wildcards.
* @param engine
* @return
*/
static std::vector<std::pair<std::string, std::string>>
compressedPaths(const engine::Engine &engine);

/**
* Sort path entries from longest to shortest
* @param paths
*/
static void sortPathsByDecreasingLength(std::vector<std::string> &paths);

/// @brief Available serializers for CAPIO-CL
struct available_serializers {
Expand All @@ -37,21 +54,23 @@ class Serializer final {
*
* @param engine instance of Engine to dump
* @param filename path of output file
* @param compress Compress the serialized output
* @throws SerializerException
*/
static void serialize_v1(const engine::Engine &engine,
const std::filesystem::path &filename);
const std::filesystem::path &filename, bool compress = false);

/**
* @brief Dump the current configuration loaded into an instance of Engine to a CAPIO-CL
* VERSION 1.1 configuration file.
*
* @param engine instance of Engine to dump
* @param filename path of output file
* @param compress Compress the serialized output
* @throws SerializerException
*/
static void serialize_v1_1(const engine::Engine &engine,
const std::filesystem::path &filename);
const std::filesystem::path &filename, bool compress = false);
};

public:
Expand All @@ -61,10 +80,11 @@ class Serializer final {
*
* @param engine instance of Engine to dump
* @param filename path of output file
* @param compress Compress directories entries when possible
* @param version Version of CAPIO-CL used to generate configuration files.
*/
static void dump(const engine::Engine &engine, const std::filesystem::path &filename,
const std::string &version = CAPIO_CL_VERSION::V1);
bool compress = false, const std::string &version = CAPIO_CL_VERSION::V1);
};
} // namespace capiocl::serializer
#endif // CAPIO_CL_SERIALIZER_H
#endif // CAPIO_CL_SERIALIZER_H
91 changes: 88 additions & 3 deletions src/Serializer.cpp
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
#include <algorithm>
#include <filesystem>
#include <fstream>
#include <jsoncons/json.hpp>
Expand All @@ -10,15 +11,15 @@

void capiocl::serializer::Serializer::dump(const engine::Engine &engine,
const std::filesystem::path &filename,
const std::string &version) {
const bool compress, const std::string &version) {
START_LOG(calf_current_tid(), "call()");
UPDATE_CALF_WORKFLOW_NAME(engine.getWorkflowName());
if (version == CAPIO_CL_VERSION::V1) {
CALF_PRINT_COLOR(CALF_CLI_LEVEL_INFO, "Serializing engine with V1 specification");
available_serializers::serialize_v1(engine, filename);
available_serializers::serialize_v1(engine, filename, compress);
} else if (version == CAPIO_CL_VERSION::V1_1) {
CALF_PRINT_COLOR(CALF_CLI_LEVEL_INFO, "Serializing engine with V1.1 specification");
available_serializers::serialize_v1_1(engine, filename);
available_serializers::serialize_v1_1(engine, filename, compress);
} else {
LOG("serializer unavailable version=%s workflow=%s output=%s", version.c_str(),
engine.getWorkflowName().c_str(), filename.string().c_str());
Expand All @@ -33,3 +34,87 @@ capiocl::serializer::SerializerException::SerializerException(const std::string
UPDATE_CALF_WORKFLOW_NAME("");
CALF_PRINT_COLOR(CALF_CLI_LEVEL_ERROR, "%s", msg.c_str());
}

void capiocl::serializer::Serializer::sortPathsByDecreasingLength(std::vector<std::string> &paths) {
std::sort(paths.begin(), paths.end(), [](const std::string &a, const std::string &b) {
if (a.length() != b.length()) {
return a.length() > b.length();
}
return a > b;
});
}

std::vector<std::pair<std::string, std::string>>
capiocl::serializer::Serializer::compressedPaths(const engine::Engine &engine) {
std::unordered_map<std::string, std::string> paths;
std::vector<std::string> directories;

for (const auto &[path, entry] : engine._capio_cl_entries) {
paths.emplace(path, path);
if (!entry.is_file || path.find_first_of("*?[") != std::string::npos) {
continue;
}

for (auto parent = std::filesystem::path(path).parent_path();;) {
if (const auto value = parent.string();
std::find(directories.begin(), directories.end(), value) == directories.end()) {
directories.push_back(value);
}
if (parent.empty() || parent == parent.root_path()) {
break;
}
parent = parent.parent_path();
}
}

sortPathsByDecreasingLength(directories);

for (const auto &directory : directories) {
const auto wildcard = (std::filesystem::path(directory) / "*").string();
if (paths.find(wildcard) != paths.end()) {
continue;
}

std::vector<std::vector<std::string>> groups;
for (const auto &[output, source] : paths) {
if (!engine._capio_cl_entries.at(source).is_file ||
(output == source && output.find_first_of("*?[") != std::string::npos)) {
continue;
}
const std::filesystem::path output_path(output);
const auto relative = output_path.lexically_relative(directory);
if ((!directory.empty() &&
(relative.empty() || *relative.begin() == std::filesystem::path(".."))) ||
(directory.empty() && output_path.is_absolute())) {
continue;
}

auto group = std::find_if(groups.begin(), groups.end(), [&](const auto &candidate) {
return engine._capio_cl_entries.at(paths.at(candidate.front())) ==
engine._capio_cl_entries.at(source);
});
(group == groups.end() ? groups.emplace_back() : *group).push_back(output);
}
for (auto &group : groups) {
std::sort(group.begin(), group.end());
}

const auto largest =
std::max_element(groups.begin(), groups.end(), [](const auto &left, const auto &right) {
return left.size() != right.size() ? left.size() < right.size()
: left.front() > right.front();
});
if (largest == groups.end() || largest->size() < 2) {
continue;
}

const auto source = paths.at(largest->front());
for (const auto &path : *largest) {
paths.erase(path);
}
paths.emplace(wildcard, source);
CALF_PRINT_COLOR(CALF_CLI_LEVEL_WARNING, "Compressing entries to %s", wildcard.c_str());
}

return {paths.begin(), paths.end()};
}
37 changes: 32 additions & 5 deletions src/serializers/v1.1.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -7,14 +7,26 @@
#include "capiocl/serializer.h"

void capiocl::serializer::Serializer::available_serializers::serialize_v1_1(
const engine::Engine &engine, const std::filesystem::path &filename) {
const engine::Engine &engine, const std::filesystem::path &filename, const bool compress) {
START_LOG(calf_current_tid(), "call()");
UPDATE_CALF_WORKFLOW_NAME(engine.getWorkflowName());

if (compress) {
CALF_PRINT_COLOR(CALF_CLI_LEVEL_WARNING, "Using configuration compression to directories!");
}

jsoncons::json doc;
doc["version"] = 1.1;
doc["name"] = engine.getWorkflowName();

const auto files = engine._capio_cl_entries;
auto files = engine._capio_cl_entries;
if (compress) {
decltype(files) compressed;
for (const auto &[path, source] : compressedPaths(engine)) {
compressed.emplace(path, files.at(source));
}
files = std::move(compressed);
}

std::unordered_map<std::string, std::vector<std::string>> app_inputs;
std::unordered_map<std::string, std::vector<std::string>> app_outputs;
Expand All @@ -27,7 +39,17 @@ void capiocl::serializer::Serializer::available_serializers::serialize_v1_1(
jsoncons::json storage = jsoncons::json::object();
jsoncons::json io_graph = jsoncons::json::array();

for (const auto &[path, entry] : files) {
std::vector<std::string> keys;
keys.reserve(files.size());
for (const auto &[k, v] : files) {
keys.push_back(k);
}

sortPathsByDecreasingLength(keys);

for (const auto &path : keys) {
const auto entry = files.at(path);

if (entry.permanent) {
permanent.push_back(path);
}
Expand All @@ -44,13 +66,18 @@ void capiocl::serializer::Serializer::available_serializers::serialize_v1_1(
}
}

for (const auto &[app_name, outputs] : app_outputs) {
for (auto &[app_name, outputs] : app_outputs) {
jsoncons::json app = jsoncons::json::object();
jsoncons::json streaming = jsoncons::json::array();
std::vector<std::string> filtered_outputs;

sortPathsByDecreasingLength(outputs);

for (const auto &path : outputs) {
const auto &entry = files.at(path);

filtered_outputs.push_back(path);

jsoncons::json streaming_item = jsoncons::json::object();
std::string committed = entry.commit_rule;
const char *name_kind = entry.is_file ? "name" : "dirname";
Expand Down Expand Up @@ -90,7 +117,7 @@ void capiocl::serializer::Serializer::available_serializers::serialize_v1_1(

app["name"] = app_name;
app["input_stream"] = app_inputs[app_name];
app["output_stream"] = outputs;
app["output_stream"] = filtered_outputs;
app["streaming"] = streaming;

io_graph.push_back(app);
Expand Down
36 changes: 31 additions & 5 deletions src/serializers/v1.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,24 @@
#include "capiocl/serializer.h"

void capiocl::serializer::Serializer::available_serializers::serialize_v1(
const engine::Engine &engine, const std::filesystem::path &filename) {
const engine::Engine &engine, const std::filesystem::path &filename, const bool compress) {
START_LOG(calf_current_tid(), "call()");
UPDATE_CALF_WORKFLOW_NAME(engine.getWorkflowName());

if (compress) {
CALF_PRINT_COLOR(CALF_CLI_LEVEL_WARNING, "Using configuration compression to directories!");
}
jsoncons::json doc;
doc["name"] = engine.getWorkflowName();

const auto files = engine._capio_cl_entries;
auto files = engine._capio_cl_entries;
if (compress) {
decltype(files) compressed;
for (const auto &[path, source] : compressedPaths(engine)) {
compressed.emplace(path, files.at(source));
}
files = std::move(compressed);
}

std::unordered_map<std::string, std::vector<std::string>> app_inputs;
std::unordered_map<std::string, std::vector<std::string>> app_outputs;
Expand All @@ -26,7 +37,17 @@ void capiocl::serializer::Serializer::available_serializers::serialize_v1(
jsoncons::json storage = jsoncons::json::object();
jsoncons::json io_graph = jsoncons::json::array();

for (const auto &[path, entry] : files) {
std::vector<std::string> keys;
keys.reserve(files.size());
for (const auto &[k, v] : files) {
keys.push_back(k);
}

sortPathsByDecreasingLength(keys);

for (const auto &path : keys) {
const auto entry = files.at(path);

if (entry.permanent) {
permanent.push_back(path);
}
Expand All @@ -43,13 +64,18 @@ void capiocl::serializer::Serializer::available_serializers::serialize_v1(
}
}

for (const auto &[app_name, outputs] : app_outputs) {
for (auto &[app_name, outputs] : app_outputs) {
jsoncons::json app = jsoncons::json::object();
jsoncons::json streaming = jsoncons::json::array();
std::vector<std::string> filtered_outputs;

sortPathsByDecreasingLength(outputs);

for (const auto &path : outputs) {
const auto &entry = files.at(path);

filtered_outputs.push_back(path);

jsoncons::json streaming_item = jsoncons::json::object();
std::string committed = entry.commit_rule;
const char *name_kind = entry.is_file ? "name" : "dirname";
Expand Down Expand Up @@ -89,7 +115,7 @@ void capiocl::serializer::Serializer::available_serializers::serialize_v1(

app["name"] = app_name;
app["input_stream"] = app_inputs[app_name];
app["output_stream"] = outputs;
app["output_stream"] = filtered_outputs;
app["streaming"] = streaming;

io_graph.push_back(app);
Expand Down
5 changes: 3 additions & 2 deletions tests/cpp/test_exceptions.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -36,8 +36,9 @@ TEST(EXCEPTION_SUITE_NAME, testFailedserializeVersion) {
const std::filesystem::path source = "/tmp/capio_cl_jsons/V" + version + "/test24.json";
auto engine = capiocl::parser::Parser::parse(source, "/tmp");

EXPECT_THROW(capiocl::serializer::Serializer::dump(*engine, "test.json", "1234.5678"),
capiocl::serializer::SerializerException);
EXPECT_THROW(
capiocl::serializer::Serializer::dump(*engine, "test.json", false, "1234.5678"),
capiocl::serializer::SerializerException);
}
}

Expand Down
Loading
Loading