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
2 changes: 2 additions & 0 deletions src/spider/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -10,13 +10,15 @@ set(SPIDER_CORE_SOURCES
set(SPIDER_CORE_HEADERS
core/Error.hpp
core/Data.hpp
core/KeyValueData.hpp
core/Task.hpp
core/TaskGraph.hpp
core/JobMetadata.hpp
io/BoostAsio.hpp
io/MsgPack.hpp
io/msgpack_message.hpp
io/Serializer.hpp
utils/TimedCache.hpp
storage/MetadataStorage.hpp
storage/DataStorage.hpp
storage/MysqlStorage.hpp
Expand Down
22 changes: 2 additions & 20 deletions src/spider/core/Data.hpp
Original file line number Diff line number Diff line change
@@ -1,43 +1,26 @@
#ifndef SPIDER_CORE_DATA_HPP
#define SPIDER_CORE_DATA_HPP

#include <optional>
#include <string>
#include <utility>
#include <vector>

#include <boost/uuid/random_generator.hpp>
#include <boost/uuid/uuid.hpp>

#include "../io/MsgPack.hpp" // IWYU pragma: keep
#include "../io/Serializer.hpp" // IWYU pragma: keep

namespace spider::core {
class Data {
public:
Data() { init_id(); }

explicit Data(std::string value) : m_value(std::move(value)) { init_id(); }

Data(boost::uuids::uuid id, std::string value) : m_id(id), m_value(std::move(value)) {}

Data(std::string key, std::string value) : m_key(std::move(key)), m_value(std::move(value)) {
init_id();
}

Data(boost::uuids::uuid id, std::string key, std::string value)
: m_id(id),
m_key(std::move(key)),
m_value(std::move(value)) {}

MSGPACK_DEFINE(m_id, m_key, m_value, m_locality, m_hard_locality);
Data(boost::uuids::uuid const id, std::string value) : m_id(id), m_value(std::move(value)) {}

static auto is_data() -> bool { return true; }

[[nodiscard]] auto get_id() const -> boost::uuids::uuid { return m_id; }

[[nodiscard]] auto get_key() const -> std::optional<std::string> { return m_key; }

[[nodiscard]] auto get_value() const -> std::string { return m_value; }

[[nodiscard]] auto get_locality() const -> std::vector<std::string> const& {
Expand All @@ -48,11 +31,10 @@ class Data {

void set_locality(std::vector<std::string> const& locality) { m_locality = locality; }

void set_hard_locality(bool hard) { m_hard_locality = hard; }
void set_hard_locality(bool const hard) { m_hard_locality = hard; }

private:
boost::uuids::uuid m_id;
std::optional<std::string> m_key;
std::string m_value;
std::vector<std::string> m_locality;
bool m_hard_locality = false;
Expand Down
30 changes: 30 additions & 0 deletions src/spider/core/KeyValueData.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
#ifndef SPIDER_CORE_KEYVALUEDATA_HPP
#define SPIDER_CORE_KEYVALUEDATA_HPP

#include <string>
#include <utility>

#include <boost/uuid/uuid.hpp>

namespace spider::core {
class KeyValueData {
public:
KeyValueData(std::string key, std::string value, boost::uuids::uuid const id)
: m_key{std::move(key)},
m_value{std::move(value)},
m_id{id} {}

[[nodiscard]] auto get_key() const -> std::string const& { return m_key; }

[[nodiscard]] auto get_value() const -> std::string const& { return m_value; }

[[nodiscard]] auto get_id() const -> boost::uuids::uuid const& { return m_id; }

private:
std::string m_key;
std::string m_value;
boost::uuids::uuid m_id;
};
} // namespace spider::core

#endif // SPIDER_CORE_KEYVALUEDATA_HPP
26 changes: 13 additions & 13 deletions src/spider/scheduler/FifoPolicy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
#include <string>
#include <vector>

#include <absl/container/flat_hash_map.h>
#include <boost/uuid/uuid.hpp>
#include <boost/uuid/uuid_io.hpp>
#include <fmt/format.h>
Expand Down Expand Up @@ -83,20 +82,24 @@ auto FifoPolicy::schedule_next(
metadata_store](core::Task const& task) -> std::chrono::system_clock::time_point {
boost::uuids::uuid const task_id = task.get_id();
boost::uuids::uuid job_id;
if (m_task_job_map.contains(task_id)) {
job_id = m_task_job_map[task_id];
std::optional<boost::uuids::uuid> const optional_job_id
= m_task_job_cache.get(task_id);
if (optional_job_id.has_value()) {
job_id = optional_job_id.value();
} else {
if (false == metadata_store->get_task_job_id(task_id, &job_id).success()) {
throw std::runtime_error(fmt::format(
"Task with id {} not exists.",
boost::uuids::to_string(task_id)
));
}
m_task_job_map.emplace(task_id, job_id);
m_task_job_cache.put(task_id, job_id);
}

if (m_job_time_map.contains(job_id)) {
return m_job_time_map[job_id];
std::optional<std::chrono::system_clock::time_point> const optional_time
= m_job_time_cache.get(job_id);
if (optional_time.has_value()) {
return optional_time.value();
}

core::JobMetadata job_metadata;
Expand All @@ -106,20 +109,17 @@ auto FifoPolicy::schedule_next(
boost::uuids::to_string(job_id)
));
}
m_job_time_map.emplace(job_id, job_metadata.get_creation_time());
m_job_time_cache.put(job_id, job_metadata.get_creation_time());
return job_metadata.get_creation_time();
}
);

return earliest_task->get_id();
}

auto FifoPolicy::cleanup_job(boost::uuids::uuid const job_id) -> void {
absl::erase_if(m_task_job_map, [&job_id](auto const& item) -> bool {
auto const& [item_task_id, item_job_id] = item;
return item_job_id == job_id;
});
m_job_time_map.erase(job_id);
auto FifoPolicy::cleanup() -> void {
m_task_job_cache.cleanup();
m_job_time_cache.cleanup();
}

} // namespace spider::scheduler
8 changes: 4 additions & 4 deletions src/spider/scheduler/FifoPolicy.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,11 @@
#include <optional>
#include <string>

#include <absl/container/flat_hash_map.h>
#include <boost/uuid/uuid.hpp>

#include "../storage/DataStorage.hpp"
#include "../storage/MetadataStorage.hpp"
#include "../utils/TimedCache.hpp"
#include "SchedulerPolicy.hpp"

namespace spider::scheduler {
Expand All @@ -23,11 +23,11 @@ class FifoPolicy final : public SchedulerPolicy {
boost::uuids::uuid worker_id,
std::string const& worker_addr
) -> std::optional<boost::uuids::uuid> override;
auto cleanup_job(boost::uuids::uuid job_id) -> void override;
auto cleanup() -> void override;

private:
absl::flat_hash_map<boost::uuids::uuid, boost::uuids::uuid> m_task_job_map;
absl::flat_hash_map<boost::uuids::uuid, std::chrono::system_clock::time_point> m_job_time_map;
core::TimedCache<boost::uuids::uuid, boost::uuids::uuid> m_task_job_cache;
core::TimedCache<boost::uuids::uuid, std::chrono::system_clock::time_point> m_job_time_cache;
};

} // namespace spider::scheduler
Expand Down
2 changes: 1 addition & 1 deletion src/spider/scheduler/SchedulerPolicy.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ class SchedulerPolicy {
std::string const& worker_addr
) -> std::optional<boost::uuids::uuid> = 0;

virtual auto cleanup_job(boost::uuids::uuid job_id) -> void = 0;
virtual auto cleanup() -> void = 0;
};

} // namespace spider::scheduler
Expand Down
19 changes: 17 additions & 2 deletions src/spider/storage/DataStorage.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@

#include "../core/Data.hpp"
#include "../core/Error.hpp"
#include "../core/KeyValueData.hpp"

namespace spider::core {
class DataStorage {
Expand All @@ -22,9 +23,9 @@ class DataStorage {
virtual void close() = 0;
virtual auto initialize() -> StorageErr = 0;

virtual auto add_data(Data const& data) -> StorageErr = 0;
virtual auto add_driver_data(boost::uuids::uuid driver_id, Data const& data) -> StorageErr = 0;
virtual auto add_task_data(boost::uuids::uuid task_id, Data const& data) -> StorageErr = 0;
virtual auto get_data(boost::uuids::uuid id, Data* data) -> StorageErr = 0;
virtual auto get_data_by_key(std::string const& key, Data* data) -> StorageErr = 0;
virtual auto remove_data(boost::uuids::uuid id) -> StorageErr = 0;
virtual auto add_task_reference(boost::uuids::uuid id, boost::uuids::uuid task_id) -> StorageErr
= 0;
Expand All @@ -34,6 +35,20 @@ class DataStorage {
add_driver_reference(boost::uuids::uuid id, boost::uuids::uuid driver_id) -> StorageErr = 0;
virtual auto
remove_driver_reference(boost::uuids::uuid id, boost::uuids::uuid driver_id) -> StorageErr = 0;
virtual auto remove_dangling_data() -> StorageErr = 0;

virtual auto add_client_kv_data(KeyValueData const& data) -> StorageErr = 0;
virtual auto add_task_kv_data(KeyValueData const& data) -> StorageErr = 0;
virtual auto get_client_kv_data(
boost::uuids::uuid const& client_id,
std::string const& key,
std::string* value
) -> StorageErr = 0;
virtual auto get_task_kv_data(
boost::uuids::uuid const& task_id,
std::string const& key,
std::string* value
) -> StorageErr = 0;
};
} // namespace spider::core

Expand Down
Loading