From c9eefc79f17f06899cf16714d46b1e6ea811df42 Mon Sep 17 00:00:00 2001 From: Marko Zivanovic Date: Sun, 6 Sep 2015 00:35:19 +0200 Subject: [PATCH] Implement file reader --- src/Facility.cpp | 4 +- src/Facility.h | 16 ++--- src/FileReader.h | 54 +++++++++++++++++ src/RFC3164FormattedSyslogMessage.h | 2 +- src/Reader.h | 2 +- src/Severity.cpp | 4 +- src/Severity.h | 14 ++--- src/SyslogMessage.cpp | 13 ++-- src/SyslogMessage.h | 2 +- src/UDPWriter.cpp | 7 ++- src/UDPWriter.h | 2 +- src/Writer.h | 2 +- test/CMakeLists.txt | 6 +- test/SyslogBulkUploaderTests.cpp | 53 +++++++++++++++-- test/UdpSyslogServer.h | 92 +++++++++++++++++++++++++++++ test/sample1 | 3 + 16 files changed, 235 insertions(+), 41 deletions(-) create mode 100644 src/FileReader.h create mode 100644 test/UdpSyslogServer.h create mode 100644 test/sample1 diff --git a/src/Facility.cpp b/src/Facility.cpp index d9b279f..e24febf 100644 --- a/src/Facility.cpp +++ b/src/Facility.cpp @@ -40,8 +40,8 @@ const std::string Facility::readFromStream(std::istream& src) { } const uint8_t Facility::readFromString(const std::string& src) { - const auto &value = _values.find(boost::algorithm::to_lower_copy(src)); - if (value != _values.end()) { + const auto& value = VALUES.find(boost::algorithm::to_lower_copy(src)); + if (value != VALUES.end()) { return value->second; } else { throw "Illegal facility value: " + src; diff --git a/src/Facility.h b/src/Facility.h index f2f8d4f..f7ff509 100644 --- a/src/Facility.h +++ b/src/Facility.h @@ -27,6 +27,8 @@ SOFTWARE. #include +class Severity; + class Facility { public: @@ -39,33 +41,31 @@ public: Facility(std::istream& source) : Facility(readFromStream(source)) { }; - Facility(const Facility & orig) : _value(orig._value) { + Facility(const Facility& orig) : _value(orig._value) { }; virtual ~Facility() { }; - bool operator!=(const Facility & right) const { + bool operator!=(const Facility& right) const { bool result = !(*this == right); // Reuse equals operator return result; } - bool operator==(const Facility & right) const { + bool operator==(const Facility& right) const { return _value == right._value; } friend std::ostream& operator<<(std::ostream& os, const Facility& obj) { - os << obj._value; + os << std::to_string(obj._value); return os; } - uint8_t as_int() const { - return _value; - } + friend const uint8_t operator+(const Facility& f, const Severity& s); private: - const std::map _values{ + const std::map VALUES{ {"kern", 0}, {"user", 1}, {"mail", 2}, diff --git a/src/FileReader.h b/src/FileReader.h new file mode 100644 index 0000000..b4fb795 --- /dev/null +++ b/src/FileReader.h @@ -0,0 +1,54 @@ +/* + The MIT License (MIT) + +Copyright (c) 2015 Marko Živanović + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. + */ + +#ifndef FILEREADER_H +#define FILEREADER_H + +#include +#include "Reader.h" + +class FileReader : public Reader { +public: + + FileReader(const std::string filename) : _stream(std::ifstream(filename)) { + if (!_stream.is_open()) { + throw "File not found: " + filename; + } + } + + virtual std::shared_ptr nextMessage() { + std::shared_ptr ret; + if (_stream.is_open() && !_stream.eof() && !_stream.fail()) { + ret.reset(new SyslogMessage(_stream)); + for (auto c = _stream.peek(); c == '\n' || c == '\r'; _stream.get(), c = _stream.peek()); + } + return ret; + } + +private: + std::ifstream _stream; +}; + +#endif /* FILEREADER_H */ + diff --git a/src/RFC3164FormattedSyslogMessage.h b/src/RFC3164FormattedSyslogMessage.h index ae94f53..20f0ccb 100644 --- a/src/RFC3164FormattedSyslogMessage.h +++ b/src/RFC3164FormattedSyslogMessage.h @@ -41,7 +41,7 @@ public: virtual ~RFC3164FormattedSyslogMessage() { }; - std::string operator()() { + const std::string operator()() { _stream.str(""); _stream << "<" << std::to_string(_message.priority()) << ">"; _stream << _message.timestamp() << " "; diff --git a/src/Reader.h b/src/Reader.h index 4092f48..9f69479 100644 --- a/src/Reader.h +++ b/src/Reader.h @@ -32,7 +32,7 @@ class SyslogMessage; class Reader : private boost::noncopyable { public: - virtual std::shared_ptr nextMessage() = 0; + virtual std::shared_ptr nextMessage() = 0; }; #endif /* READER_H */ diff --git a/src/Severity.cpp b/src/Severity.cpp index 881f155..3cc700a 100644 --- a/src/Severity.cpp +++ b/src/Severity.cpp @@ -40,8 +40,8 @@ const std::string Severity::readFromStream(std::istream& src) { } const uint8_t Severity::readFromString(const std::string& src) { - const auto &value = _values.find(boost::algorithm::to_lower_copy(src)); - if (value != _values.end()) { + const auto &value = VALUES.find(boost::algorithm::to_lower_copy(src)); + if (value != VALUES.end()) { return value->second; } else { throw "illegal severity: " + src; diff --git a/src/Severity.h b/src/Severity.h index cdca9dc..fad1d0a 100644 --- a/src/Severity.h +++ b/src/Severity.h @@ -25,6 +25,8 @@ SOFTWARE. #ifndef SEVERITY_H #define SEVERITY_H +class Facility; + class Severity { public: @@ -52,18 +54,16 @@ public: return _value == right._value; } - friend std::ostream& operator<<(std::ostream& os, const Severity& obj) { - os << obj._value; - return os; - } + friend const uint8_t operator+(const Facility& f, const Severity& s); - uint8_t as_int() const { - return _value; + friend std::ostream& operator<<(std::ostream& os, const Severity& obj) { + os << std::to_string(obj._value); + return os; } private: - const std::map _values{ + const std::map VALUES{ {"emergency", 0}, {"alert", 1}, {"critical", 2}, diff --git a/src/SyslogMessage.cpp b/src/SyslogMessage.cpp index 396f283..0e35a7e 100644 --- a/src/SyslogMessage.cpp +++ b/src/SyslogMessage.cpp @@ -76,15 +76,14 @@ const std::string SyslogMessage::readSource(std::istream& src) { return ret; }; -const std::string SyslogMessage::readMessage(std::istream & src) { +const std::string SyslogMessage::readMessage(std::istream& src) { std::string ret; skipWhitespace(src); - while (src) { - auto c = src.get(); - if (!src.eof()) { - ret.push_back(c); - } - } + std::getline(src, ret); return ret; }; +const uint8_t operator+(const Facility& f, const Severity& s) { + return f._value * 8 + s._value; +}; + diff --git a/src/SyslogMessage.h b/src/SyslogMessage.h index 9f08779..c8d0708 100644 --- a/src/SyslogMessage.h +++ b/src/SyslogMessage.h @@ -65,7 +65,7 @@ public: } const uint8_t priority() const { - return (_facility.as_int() * 8) +_severity.as_int(); + return _facility + _severity; }; friend std::ostream& operator<<(std::ostream& os, const SyslogMessage& obj) { diff --git a/src/UDPWriter.cpp b/src/UDPWriter.cpp index 28d304f..ce7abe5 100644 --- a/src/UDPWriter.cpp +++ b/src/UDPWriter.cpp @@ -26,7 +26,8 @@ SOFTWARE. #include "UDPWriter.h" #include "RFC3164FormattedSyslogMessage.h" -using namespace boost::asio::ip; +using boost::asio::ip::udp; +using namespace boost::asio; UDPWriter::UDPWriter(const std::string& destination, const int port) { udp::resolver resolver(_ios); @@ -40,10 +41,10 @@ UDPWriter::UDPWriter(const std::string& destination, const int port) { } }; -void UDPWriter::sendMessage(std::shared_ptr message) { +void UDPWriter::sendMessage(std::shared_ptr message) { if (_socket.is_open()) { auto formatted = RFC3164FormattedSyslogMessage{*message}; - auto buff = ::boost::asio::buffer(formatted()); + auto buff = buffer(formatted()); _socket.send(buff); } } diff --git a/src/UDPWriter.h b/src/UDPWriter.h index ce21366..3a85e28 100644 --- a/src/UDPWriter.h +++ b/src/UDPWriter.h @@ -34,7 +34,7 @@ public: UDPWriter(const std::string&, const int); - virtual void sendMessage(std::shared_ptr); + virtual void sendMessage(std::shared_ptr); private: boost::asio::io_service _ios; diff --git a/src/Writer.h b/src/Writer.h index 40000ce..16aaf15 100644 --- a/src/Writer.h +++ b/src/Writer.h @@ -31,7 +31,7 @@ SOFTWARE. class Writer : private boost::noncopyable { public: - virtual void sendMessage(std::shared_ptr) = 0; + virtual void sendMessage(std::shared_ptr) = 0; }; #endif /* WRITER_H */ diff --git a/test/CMakeLists.txt b/test/CMakeLists.txt index cb4be34..98b07ee 100644 --- a/test/CMakeLists.txt +++ b/test/CMakeLists.txt @@ -1,4 +1,5 @@ -find_package(Boost COMPONENTS unit_test_framework date_time REQUIRED) +find_package(Boost COMPONENTS unit_test_framework date_time filesystem system REQUIRED) +find_package(Threads) include_directories( ${TEST_SOURCE_DIR/src} ${Boost_INLUDE_DIRS} @@ -10,6 +11,9 @@ target_link_libraries(SyslogBulkUploaderTests slbu-lib ${Boost_UNIT_TEST_FRAMEWORK_LIBRARY} ${Boost_DATE_TIME_LIBRARY} + ${Boost_FILESYSTEM_LIBRARY} + ${Boost_SYSTEM_LIBRARY} + ${CMAKE_THREAD_LIBS_INIT} ) add_executable(SyslogMessageTests SyslogMessageTests.cpp) diff --git a/test/SyslogBulkUploaderTests.cpp b/test/SyslogBulkUploaderTests.cpp index f0b3c52..efe187a 100644 --- a/test/SyslogBulkUploaderTests.cpp +++ b/test/SyslogBulkUploaderTests.cpp @@ -25,14 +25,21 @@ SOFTWARE. #include "../src/SyslogBulkUploader.h" #include "../src/Reader.h" #include "../src/Writer.h" +#include "../src/FileReader.h" +#include "../src/UDPWriter.h" +#include "UdpSyslogServer.h" +#include "../src/RFC3164FormattedSyslogMessage.h" #define BOOST_TEST_MODULE SyslogBulkUploaderTests #include #include +#define BOOST_FILESYSTEM_NO_DEPRECATED +#include +#include class MockReader : public Reader { public: - virtual std::shared_ptr nextMessage() { + virtual std::shared_ptr nextMessage() { if (_pos >= _messages.size()) { return std::shared_ptr(); } else { @@ -52,10 +59,11 @@ private: class MockWriter : public Writer { public: - virtual void sendMessage(std::shared_ptr msg) { + virtual void sendMessage(std::shared_ptr msg) { _messages.push_back(msg); + std::cout << *msg << std::endl; }; - std::vector> _messages; + std::vector> _messages; }; BOOST_AUTO_TEST_CASE(test_run) { @@ -64,7 +72,40 @@ BOOST_AUTO_TEST_CASE(test_run) { SyslogBulkUploader ul(r, w); ul.run(); BOOST_CHECK_EQUAL(w._messages.size(), 3); - for (auto i = 0; i < w._messages.size(); i++) { - std::cout << *(w._messages[i].get()) << std::endl; - } } + +BOOST_AUTO_TEST_CASE(test_file_reader) { + std::cout << "Running in " << boost::filesystem::initial_path() << std::endl; + FileReader r("../test/sample1"); + MockWriter w; + SyslogBulkUploader ul(r, w); + ul.run(); + BOOST_CHECK_EQUAL(w._messages.size(), 3); +} + +BOOST_AUTO_TEST_CASE(test_udp_writer) { + using namespace boost::posix_time; + + ptime start(second_clock::local_time()); + boost::asio::io_service ios; + FileReader r("../test/sample1"); + UDPWriter w("localhost", 51514); + SyslogBulkUploader ul(r, w); + UdpSyslogServer server(ios, 51514, 2000); + ul.run(); + + while (second_clock::local_time() - start < seconds(2)) { + ios.run_one(); + } + + BOOST_CHECK_EQUAL(server.messages().size(), 3); + + auto it = server.messages().begin(); + FileReader r2("../test/sample1"); + while (auto msg = r2.nextMessage()) { + RFC3164FormattedSyslogMessage fmt(*msg); + BOOST_CHECK_EQUAL(*it, fmt()); + it++; + } + +} \ No newline at end of file diff --git a/test/UdpSyslogServer.h b/test/UdpSyslogServer.h new file mode 100644 index 0000000..7272e4d --- /dev/null +++ b/test/UdpSyslogServer.h @@ -0,0 +1,92 @@ +/* + The MIT License (MIT) + +Copyright (c) 2015 Marko Živanović + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. + */ + +#ifndef UDPSYSLOGSERVER_H +#define UDPSYSLOGSERVER_H + +#include +#include +#include +#include +#include +#include + +using boost::asio::ip::udp; +using namespace boost::asio; +using namespace boost::posix_time; + +class UdpSyslogServer { +public: + + UdpSyslogServer(boost::asio::io_service& ios, const int port, const int timeout) + : _socket(ios, udp::endpoint(udp::v4(), port)), _dt(ios), _timeout(timeout) { + deadlineHandler(); + receive(); + } + + const std::vector& messages() const { + return _messages; + } + +private: + + void receive() { + _dt.expires_from_now(milliseconds(_timeout)); + _socket.async_receive_from( + boost::asio::buffer(_buffer), _peer, + boost::bind(&UdpSyslogServer::receiveHandler, this, + boost::asio::placeholders::error, + boost::asio::placeholders::bytes_transferred)); + } + + void receiveHandler(const boost::system::error_code& error, std::size_t len) { + if (!error || error == boost::asio::error::message_size) { + if (len > 0) { + std::string s(_buffer.data(), len); + _messages.push_back(s); + } + receive(); + } + } + + void deadlineHandler() { + if (_dt.expires_at() <= deadline_timer::traits_type::now()) { + _socket.cancel(); + _dt.expires_at(boost::posix_time::pos_infin); + } + + _dt.async_wait(boost::bind(&UdpSyslogServer::deadlineHandler, this)); + } + + + udp::socket _socket; + udp::endpoint _peer; + boost::array _buffer; + std::vector _messages; + deadline_timer _dt; + const int _timeout; +}; + +#endif /* UDPSYSLOGSERVER_H */ + diff --git a/test/sample1 b/test/sample1 new file mode 100644 index 0000000..a8a7b5b --- /dev/null +++ b/test/sample1 @@ -0,0 +1,3 @@ +2015-09-02 13:33:11 Local4.Critical 192.168.0.1 Kiwi_Syslog_Server %ASA-2-106007: Deny inbound UDP from 1.2.3.4/22084 to 4.3.2.1/53 due to DNS Query +2015-09-02 13:33:11 Local4.Critical 192.168.0.1 Kiwi_Syslog_Server %ASA-2-106007: Deny inbound UDP from 1.2.3.4/22084 to 4.3.2.1/53 due to DNS Query +2015-09-02 13:33:11 Local4.Critical 192.168.0.1 Kiwi_Syslog_Server %ASA-2-106007: Deny inbound UDP from 1.2.3.4/22084 to 4.3.2.1/53 due to DNS Query