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
436 changes: 436 additions & 0 deletions broker/bam/doc/bam.md

Large diffs are not rendered by default.

19 changes: 11 additions & 8 deletions broker/bam/test/monitoring_stream.cc
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,9 @@ using log_v2 = com::centreon::common::log_v2::log_v2;
using namespace com::centreon::broker;
using namespace com::centreon::broker::bam;

const std::string db_user = "root";
const std::string db_password = "centreon";

class BamMonitoringStream : public testing::Test {
void SetUp() override {
config::applier::init(com::centreon::common::BROKER, 0, "test_broker", 0);
Expand All @@ -38,9 +41,9 @@ class BamMonitoringStream : public testing::Test {
};

TEST_F(BamMonitoringStream, WriteKpi) {
database_config cfg("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config cfg("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon");
database_config storage("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config storage("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon_storage");

std::shared_ptr<persistent_cache> cache;
Expand All @@ -56,9 +59,9 @@ TEST_F(BamMonitoringStream, WriteKpi) {
}

TEST_F(BamMonitoringStream, WriteBA) {
database_config cfg("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config cfg("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon");
database_config storage("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config storage("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon_storage");
;
std::shared_ptr<persistent_cache> cache;
Expand All @@ -73,9 +76,9 @@ TEST_F(BamMonitoringStream, WriteBA) {
}

TEST_F(BamMonitoringStream, WorkWithNoPendigMysqlRequest) {
database_config cfg("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config cfg("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon", 0);
database_config storage("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config storage("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon_storage", 0);
;
std::shared_ptr<persistent_cache> cache;
Expand All @@ -95,9 +98,9 @@ TEST_F(BamMonitoringStream, WorkWithNoPendigMysqlRequest) {
}

TEST_F(BamMonitoringStream, WorkWithPendigMysqlRequest) {
database_config cfg("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config cfg("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon", 5);
database_config storage("MySQL", "127.0.0.1", "", 3306, "root", "centreon",
database_config storage("MySQL", "127.0.0.1", "", 3306, db_user, db_password,
"centreon_storage", 5);
;
std::shared_ptr<persistent_cache> cache;
Expand Down
29 changes: 22 additions & 7 deletions broker/core/sql/src/query_preparator.cc
Original file line number Diff line number Diff line change
Expand Up @@ -548,20 +548,35 @@ mysql_stmt query_preparator::prepare_update_table(
query.append("=?,");
query_bind_mapping.insert(std::make_pair(key, query_size++));
}
// Part of ID field.
else {
where.append(e.name);
where.append("=? AND ");
key = fmt::format(":{}", entry_name);
where_bind_mapping.insert(std::make_pair(key, where_size++));
}
} else
throw msg_fmt(
"could not prepare update query for event of type {}:"
"protobuf field with number {} does not exist in '{}' protobuf "
"object",
_event_id, e.number, info->get_name());
}

for (const auto& e : _pb_unique) {
const google::protobuf::FieldDescriptor* f =
desc->FindFieldByNumber(e.number);
if (!f)
throw msg_fmt(
"could not prepare update query for event of type {}:"
"protobuf field with number {} does not exist in '{}' protobuf "
"object",
_event_id, e.number, info->get_name());
std::string_view entry_name = f->name();
if (static_cast<uint32_t>(f->index()) >= pb_mapping.size())
pb_mapping.resize(f->index() + 1);
if (std::get<0>(pb_mapping[f->index()]).empty())
pb_mapping[f->index()] = std::make_tuple(entry_name, e.max_length, e.attribute);

where.append(e.name);
where.append("=? AND ");
key = fmt::format(":{}", entry_name);
where_bind_mapping.insert(std::make_pair(key, where_size++));
}

query.resize(query.size() - 1);
query.append(where, 0, where.size() - 5);

Expand Down
75 changes: 39 additions & 36 deletions broker/core/test/mysql/mysql.cc
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,9 @@ using namespace com::centreon::broker;
using namespace com::centreon::broker::database;
using log_v2 = com::centreon::common::log_v2::log_v2;

const std::string db_user = "root";
const std::string db_password = "centreon";

class DatabaseStorageTest : public ::testing::Test {
public:
void SetUp() override {
Expand All @@ -65,8 +68,8 @@ class DatabaseStorageTest : public ::testing::Test {
// When there is no database
// Then the mysql creation throws an exception
TEST_F(DatabaseStorageTest, NoDatabase) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 9876, "root",
"centreon", "centreon_storage");
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 9876, db_user,
db_password, "centreon_storage");
std::unique_ptr<mysql> ms;
ASSERT_THROW(ms.reset(new mysql(db_cfg)), msg_fmt);
}
Expand All @@ -75,8 +78,8 @@ TEST_F(DatabaseStorageTest, NoDatabase) {
// And when the connection is well done
// Then no exception is thrown and the mysql object is well built.
TEST_F(DatabaseStorageTest, ConnectionOk) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage");
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage");
std::unique_ptr<mysql> ms;
ASSERT_NO_THROW(ms = std::make_unique<mysql>(db_cfg));
}
Expand Down Expand Up @@ -506,8 +509,8 @@ TEST_F(DatabaseStorageTest, ConnectionOk) {
TEST_F(DatabaseStorageTest, CustomVarStatement) {
config::applier::modules modules(log_v2::instance().get(log_v2::SQL));
modules.load_file("./broker/lib/10-neb.so");
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
std::unique_ptr<mysql> ms(new mysql(db_cfg));
query_preparator::event_unique unique;
unique.insert("host_id");
Expand Down Expand Up @@ -1258,8 +1261,8 @@ TEST_F(DatabaseStorageTest, CustomVarStatement) {
////}
//
TEST_F(DatabaseStorageTest, ChooseConnectionByName) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms = std::make_unique<mysql>(db_cfg);
int thread_foo(ms->choose_connection_by_name("foo"));
int thread_bar(ms->choose_connection_by_name("bar"));
Expand All @@ -1280,8 +1283,8 @@ TEST_F(DatabaseStorageTest, ChooseConnectionByName) {
// Then we can bind values to it and execute the statement.
// Then a commit makes data available in the database.
TEST_F(DatabaseStorageTest, RepeatStatements) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -1359,8 +1362,8 @@ TEST_F(DatabaseStorageTest, RepeatStatements) {
}

TEST_F(DatabaseStorageTest, CheckBulkStatement) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string version = ms->get_server_version();
std::vector<std::string_view> arr =
Expand Down Expand Up @@ -1428,8 +1431,8 @@ TEST_F(DatabaseStorageTest, CheckBulkStatement) {
TEST_F(DatabaseStorageTest, UpdateBulkStatement) {
constexpr int TOTAL = 20;

database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
if (ms->support_bulk_statement()) {
std::string query{
Expand Down Expand Up @@ -1492,8 +1495,8 @@ TEST_F(DatabaseStorageTest, UpdateBulkStatement) {
// Then we can bind values to it and execute the statement.
// Then a commit makes data available in the database.
TEST_F(DatabaseStorageTest, LastInsertId) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
time_t now = time(nullptr);
std::string query(
fmt::format("INSERT INTO metrics"
Expand Down Expand Up @@ -1533,8 +1536,8 @@ TEST_F(DatabaseStorageTest, LastInsertId) {
}

TEST_F(DatabaseStorageTest, BulkStatementWithNullStr) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
if (ms->support_bulk_statement()) {
std::string query1{"DROP TABLE IF EXISTS ut_test"};
Expand Down Expand Up @@ -1596,8 +1599,8 @@ TEST_F(DatabaseStorageTest, BulkStatementWithNullStr) {
}

TEST_F(DatabaseStorageTest, RepeatStatementsWithNull) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -1645,8 +1648,8 @@ TEST_F(DatabaseStorageTest, RepeatStatementsWithNull) {
}

TEST_F(DatabaseStorageTest, RepeatStatementsWithBigStrings) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -1747,8 +1750,8 @@ TEST_F(DatabaseStorageTest, RepeatStatementsWithBigStrings) {
}

TEST_F(DatabaseStorageTest, RepeatStatementsWithNullValues) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -1825,8 +1828,8 @@ TEST_F(DatabaseStorageTest, RepeatStatementsWithNullValues) {
}

TEST_F(DatabaseStorageTest, BulkStatementsWithNullValues) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -1933,8 +1936,8 @@ TEST_F(DatabaseStorageTest, BulkStatementsWithNullValues) {
}

TEST_F(DatabaseStorageTest, RepeatStatementsWithBooleanValues) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -1973,8 +1976,8 @@ TEST_F(DatabaseStorageTest, RepeatStatementsWithBooleanValues) {
}

TEST_F(DatabaseStorageTest, BulkStatementsWithBooleanValues) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -2042,8 +2045,8 @@ static std::string row_filler2(const row& data) {
}

TEST_F(DatabaseStorageTest, MySqlMultiInsert) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -2191,8 +2194,8 @@ struct multi_event_binder {
};

TEST_F(DatabaseStorageTest, bulk_or_multi_bbdo_event_bulk) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down Expand Up @@ -2246,8 +2249,8 @@ TEST_F(DatabaseStorageTest, bulk_or_multi_bbdo_event_bulk) {
}

TEST_F(DatabaseStorageTest, bulk_or_multi_bbdo_event_multi) {
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, "root",
"centreon", "centreon_storage", 5, true, 5);
database_config db_cfg("MySQL", "127.0.0.1", MYSQL_SOCKET, 3306, db_user,
db_password, "centreon_storage", 5, true, 5);
auto ms{std::make_unique<mysql>(db_cfg)};
std::string query1{"DROP TABLE IF EXISTS ut_test"};
std::string query2{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -261,8 +261,9 @@ class stream : public io::stream {
std::shared_ptr<stats::center> _center;
ConflictManagerStats* _stats;

absl::flat_hash_set<uint32_t> _cache_deleted_instance_id;
std::unordered_map<uint32_t, uint32_t> _cache_host_instance;
absl::flat_hash_set<uint64_t> _cache_deleted_instance_id;
std::unordered_map<uint64_t /*host_id*/, uint64_t /*instance_id*/>
_cache_host_instance;
absl::flat_hash_map<uint64_t, size_t> _cache_hst_cmd;
absl::flat_hash_map<std::pair<uint64_t, uint64_t>, size_t> _cache_svc_cmd;
absl::flat_hash_map<std::pair<uint64_t, uint64_t>, index_info> _index_cache;
Expand All @@ -274,7 +275,11 @@ class stream : public io::stream {
absl::flat_hash_map<std::pair<uint64_t, uint16_t>, uint64_t> _severity_cache;
absl::flat_hash_map<std::pair<uint64_t, uint16_t>, uint64_t> _tags_cache;

absl::flat_hash_map<std::pair<uint64_t, uint64_t>, uint64_t> _resource_cache;
absl::flat_hash_map<std::tuple<uint64_t /*id*/,
uint64_t /*parent_id*/,
uint64_t /*instance_id*/>,
uint64_t>
_resource_cache;

mutable absl::Mutex _timer_m;
/* This is a barrier for timers. It must be locked in shared mode in the
Expand Down Expand Up @@ -329,11 +334,13 @@ class stream : public io::stream {
database::mysql_stmt _pb_host_check_update;
database::mysql_stmt _host_group_insupdate;
database::mysql_stmt _pb_host_group_insupdate;
database::mysql_stmt _host_group_member_delete;
std::unique_ptr<database::mysql_stmt_base> _host_group_member_delete;
database::mysql_stmt _host_group_member_insert;
database::mysql_stmt _pb_host_group_member_insert;
database::mysql_stmt _host_insupdate;
database::mysql_stmt _host_update;
database::mysql_stmt _pb_host_insupdate;
database::mysql_stmt _pb_host_update;
database::mysql_stmt _host_parent_delete;
database::mysql_stmt _host_parent_insert;
database::mysql_stmt _pb_host_parent_delete;
Expand All @@ -347,12 +354,12 @@ class stream : public io::stream {
database::mysql_stmt _pb_service_check_update;
database::mysql_stmt _service_group_insupdate;
database::mysql_stmt _pb_service_group_insupdate;
database::mysql_stmt _service_group_member_delete;
std::unique_ptr<database::mysql_stmt_base> _service_group_member_delete;
database::mysql_stmt _service_group_member_insert;
database::mysql_stmt _pb_service_group_member_delete;
database::mysql_stmt _pb_service_group_member_insert;
database::mysql_stmt _service_insupdate;
database::mysql_stmt _pb_service_insupdate;
database::mysql_stmt _pb_service_update;
database::mysql_stmt _service_status_update;

std::unique_ptr<database::mysql_stmt_base> _hscr_update;
Expand Down Expand Up @@ -451,6 +458,7 @@ class stream : public io::stream {
void _process_service_status(const std::shared_ptr<io::data>& d);
void _process_responsive_instance(const std::shared_ptr<io::data>& d);

void _prepare_pb_requests();
void _process_pb_host(const std::shared_ptr<io::data>& d);
uint64_t _process_pb_host_in_resources(const Host& h, int32_t conn);
void _process_pb_instance_configuration(const std::shared_ptr<io::data>& d);
Expand Down
Loading