diff --git a/osquery/dispatcher/distributed_runner.cpp b/osquery/dispatcher/distributed_runner.cpp index 39ff56bb95a..75773a20b7d 100644 --- a/osquery/dispatcher/distributed_runner.cpp +++ b/osquery/dispatcher/distributed_runner.cpp @@ -8,11 +8,13 @@ */ #include +#include #include #include #include #include +#include #include #include @@ -26,15 +28,42 @@ FLAG(uint64, "Seconds between polling for new queries (default 60)") DECLARE_bool(disable_distributed); +DECLARE_bool(openframe_mode); DECLARE_string(distributed_plugin); const size_t kDistributedAccelerationInterval = 5; +static void logOutcomeChange(const char* operation, + const Status& status, + std::string& last) { + const auto message = status.ok() ? std::string() : status.getMessage(); + if (message == last) { + return; + } + + last = message; + if (status.ok()) { + LOG(INFO) << operation << ": recovered"; + } else { + LOG(ERROR) << operation << ": " << message; + } +} + void DistributedRunner::start() { auto dist = Distributed(); + std::string last_read_error; + std::string last_write_error; + while (!interrupted()) { - dist.pullUpdates(); - dist.runQueries(); + const auto read_status = dist.pullUpdates(); + const auto write_status = dist.runQueries(); + if (FLAGS_openframe_mode) { + logOutcomeChange( + "Reading distributed queries", read_status, last_read_error); + logOutcomeChange( + "Writing distributed query results", write_status, last_write_error); + } + dist.cleanupExpiredRunningQueries(); std::string accelerate_checkins_expire_str = "-1"; diff --git a/osquery/distributed/distributed.cpp b/osquery/distributed/distributed.cpp index e9d313a2122..cf3f24d6fa4 100644 --- a/osquery/distributed/distributed.cpp +++ b/osquery/distributed/distributed.cpp @@ -49,6 +49,7 @@ FLAG(uint64, "Seconds to denylist distributed queries (default 1 day)"); DECLARE_bool(verbose); +DECLARE_bool(openframe_mode); std::string Distributed::currentRequestId_{""}; @@ -153,6 +154,7 @@ void Distributed::addResult(const DistributedQueryResult& result) { Status Distributed::runQueries() { auto queries = getPendingQueries(); + size_t failed = 0; for (const auto& query : queries) { auto request = popRequest(query); @@ -166,6 +168,7 @@ Status Distributed::runQueries() { result.status = Status(1, "Denylisted"); result.message = "distributed query is denylisted"; addResult(result); + ++failed; continue; } @@ -184,6 +187,7 @@ Status Distributed::runQueries() { const auto ok = sql.getStatus().ok(); const auto& msg = ok ? "" : sql.getMessageString(); if (!ok) { + ++failed; LOG(ERROR) << "Error executing distributed query: " << request.id << ": " << msg; } @@ -194,6 +198,12 @@ Status Distributed::runQueries() { request, sql.rows(), sql.columns(), sql.getStatus(), msg); addResult(result); } + + if (FLAGS_openframe_mode && !queries.empty()) { + LOG(INFO) << "Executed " << queries.size() << " distributed queries (" + << failed << " failed)"; + } + return flushCompleted(); } @@ -273,6 +283,10 @@ Status Distributed::flushCompleted() { {{"action", "writeResults"}, {"results", results}}, response); if (s.ok()) { + if (FLAGS_openframe_mode) { + LOG(INFO) << "Sent results for " << getCompletedCount() + << " distributed queries"; + } results_.clear(); performance_.clear(); }