From 690fac67e05f99614d372970c9de16ecfb2ada9e Mon Sep 17 00:00:00 2001 From: emeric Date: Thu, 15 May 2014 14:06:09 +0200 Subject: [PATCH] WIP, Remote Client/Server... --- Makefile.am | 6 + database/Artist.cpp | 6 + database/AudioTypes.hpp | 3 + remote/messages/Header.hpp | 39 +++-- .../server/AudioCollectionRequestHandler.cpp | 81 ++++++++++ .../server/AudioCollectionRequestHandler.hpp | 30 ++++ remote/server/Connection.cpp | 138 +++++++++++++++--- remote/server/Connection.hpp | 22 +-- remote/server/RequestHandler.cpp | 31 +++- remote/server/RequestHandler.hpp | 9 +- test/Makefile.am | 1 + test/Makefile.in | 17 +++ test/RemoteClientServer.cpp | 81 +++++++--- test/TestDatabase.cpp | 2 +- 14 files changed, 395 insertions(+), 71 deletions(-) create mode 100644 remote/server/AudioCollectionRequestHandler.cpp create mode 100644 remote/server/AudioCollectionRequestHandler.hpp diff --git a/Makefile.am b/Makefile.am index 5b38865c..93f7965b 100644 --- a/Makefile.am +++ b/Makefile.am @@ -45,8 +45,14 @@ lms_SOURCES = \ $(top_srcdir)/transcode/InputMediaFile.cpp \ $(top_srcdir)/metadata/AvFormat.cpp \ $(top_srcdir)/main/RemoteServerService.cpp \ + $(top_srcdir)/remote/messages/auth.pb.cc \ + $(top_srcdir)/remote/messages/collection.pb.cc \ + $(top_srcdir)/remote/messages/common.pb.cc \ + $(top_srcdir)/remote/messages/media.pb.cc \ + $(top_srcdir)/remote/messages/messages.pb.cc \ $(top_srcdir)/remote/server/Connection.cpp \ $(top_srcdir)/remote/server/ConnectionManager.cpp \ + $(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp \ $(top_srcdir)/remote/server/RequestHandler.cpp \ $(top_srcdir)/remote/server/Server.cpp \ $(top_srcdir)/metadata/Utils.cpp diff --git a/database/Artist.cpp b/database/Artist.cpp index 1b61539d..5bf6c3e4 100644 --- a/database/Artist.cpp +++ b/database/Artist.cpp @@ -29,3 +29,9 @@ Artist::getNone(Wt::Dbo::Session& session) return res; } + +Wt::Dbo::collection +Artist::getAll(Wt::Dbo::Session& session) +{ + return session.find(); +} diff --git a/database/AudioTypes.hpp b/database/AudioTypes.hpp index 128acb9b..d314e61a 100644 --- a/database/AudioTypes.hpp +++ b/database/AudioTypes.hpp @@ -27,6 +27,9 @@ class Artist // Accessors static pointer getByName(Wt::Dbo::Session& session, const std::string& name); static pointer getNone(Wt::Dbo::Session& session); + static Wt::Dbo::collection getAll(Wt::Dbo::Session& session); + + const std::string& getName(void) const { return _name; } // Create static pointer create(Wt::Dbo::Session& session, const std::string& name); diff --git a/remote/messages/Header.hpp b/remote/messages/Header.hpp index d6fa4709..1e1d0dc9 100644 --- a/remote/messages/Header.hpp +++ b/remote/messages/Header.hpp @@ -16,22 +16,38 @@ class Header void setSize(std::size_t size) { _size = size; } std::size_t getSize(void) const {return _size;} - static Header from_buffer(const std::array& buffer, bool& error) + + bool from_istream(std::istream &is) { - Header res; + std::array buffer; - if (decode32(&buffer[0]) != _magic) - error = true; + bool res = is.read(reinterpret_cast(buffer.data()), buffer.size()); + + if (res && is.gcount() == buffer.size()) + return from_buffer(buffer); else - { - res._size = decode32(&buffer[4]); - error = false; - } - - return res; + return false; } - void to_buffer(std::array& buffer) + bool from_buffer(const std::array& buffer) + { + if (decode32(&buffer[0]) != _magic) + return false; + else + { + _size = decode32(&buffer[4]); + + return _size < _maxSize; + } + } + + void to_ostream(std::ostream& os) + { + + + } + + void to_buffer(std::array& buffer) const { encode32(_magic, &buffer[0]); encode32(_size, &buffer[4]); @@ -55,6 +71,7 @@ class Header } static const uint32_t _magic = 0xbeef; + static const uint32_t _maxSize = 65536; uint32_t _size; }; diff --git a/remote/server/AudioCollectionRequestHandler.cpp b/remote/server/AudioCollectionRequestHandler.cpp new file mode 100644 index 00000000..2696ad2a --- /dev/null +++ b/remote/server/AudioCollectionRequestHandler.cpp @@ -0,0 +1,81 @@ +#include "AudioCollectionRequestHandler.hpp" + +#include "database/AudioTypes.hpp" + +namespace Remote { +namespace Server { + +AudioCollectionRequestHandler::AudioCollectionRequestHandler(DatabaseHandler& db) +: _db(db) +{} + + +bool +AudioCollectionRequestHandler::process(const AudioCollectionRequest& request, std::vector responses) +{ + + switch (request.type()) + { + case AudioCollectionRequest_Type_TypeGetArtistList: + if (request.has_get_artists()) + return processGetArtists(request.get_artists(), responses); + else + std::cerr << "Bad AudioCollectionRequest_Type_TypeGetArtistList: message!" << std::endl; + break; + + case AudioCollectionRequest_Type_TypeGetReleaseList: + break; + + case AudioCollectionRequest_Type_TypeGetTrackList: + break; + + default: + std::cerr << "Unhandled AudioCollectionRequest_Type = " << request.type() << std::endl; + } + + return false; +} + +bool +AudioCollectionRequestHandler::processGetArtists(const AudioCollectionRequest::GetArtistList& request, std::vector responses) +{ + + // sanity checks + if (request.has_filter_name()) + std::cout << "Filter nameĀ = " << request.filter_name() << std::endl; + + for (int id = 0; id < request.filter_genre_size(); ++id) + { + std::cout << "Filter genre " << id << " = '" << request.filter_genre(id) << "'" << std::endl; + } + + if (request.has_preferred_batch_size()) + std::cout << "Requested batch size = " << request.preferred_batch_size() << std::endl; + + + // Now fetch requested data... + + std::cout << "Getting artists..." << std::endl; + + Wt::Dbo::Transaction transaction( _db.getSession() ); + + Wt::Dbo::collection artists = Artist::getAll( _db.getSession() ); + + std::cout << "size = " << artists.size() << std::endl; + + typedef Wt::Dbo::collection< Artist::pointer > Artists; + + for (Artists::iterator it = artists.begin(); it != artists.end(); ++it) + { + std::cout << "Spotted artist = " << (*it)->getName() << std::endl; + } + + std::cout << "Getting artists DONE" << std::endl; + + return false; +} + +} // namespace Remote +} // namespace Server + + diff --git a/remote/server/AudioCollectionRequestHandler.hpp b/remote/server/AudioCollectionRequestHandler.hpp new file mode 100644 index 00000000..7d7563d3 --- /dev/null +++ b/remote/server/AudioCollectionRequestHandler.hpp @@ -0,0 +1,30 @@ +#ifndef REMOTE_AUDIO_COLLECTION_REQUEST_HANDLER +#define REMOTE_AUDIO_COLLECTION_REQUEST_HANDLER + +#include "messages/messages.pb.h" + +#include "database/DatabaseHandler.hpp" + +namespace Remote { +namespace Server { + +class AudioCollectionRequestHandler +{ + public: + AudioCollectionRequestHandler(DatabaseHandler& db); + + bool process(const AudioCollectionRequest& request, std::vector responses); + + private: + + bool processGetArtists(const AudioCollectionRequest::GetArtistList& request, std::vector responses); + + DatabaseHandler& _db; + +}; + +} // namespace Remote +} // namespace Server + +#endif + diff --git a/remote/server/Connection.cpp b/remote/server/Connection.cpp index b5332e81..2ad5da87 100644 --- a/remote/server/Connection.cpp +++ b/remote/server/Connection.cpp @@ -1,7 +1,11 @@ #include #include + #include +#include + +#include "messages/messages.pb.h" #include "RequestHandler.hpp" @@ -18,6 +22,7 @@ Connection::Connection(boost::asio::ip::tcp::socket socket, _connectionManager(manager), _requestHandler(handler) { + std::cout << "Server::Connection::Connection, Creating connection" << std::endl; } boost::asio::ip::tcp::socket& @@ -29,29 +34,59 @@ Connection::socket() void Connection::start() { - _socket.async_read_some(boost::asio::buffer(_headerBuffer), - boost::bind(&Connection::handleRead, shared_from_this(), - boost::asio::placeholders::error, - boost::asio::placeholders::bytes_transferred)); + boost::asio::streambuf::mutable_buffers_type bufs = _inputStreamBuf.prepare(Remote::Header::size); + + boost::asio::async_read(_socket, + bufs, + boost::asio::transfer_exactly(Remote::Header::size), + boost::bind(&Connection::handleReadHeader, shared_from_this(), + boost::asio::placeholders::error, + boost::asio::placeholders::bytes_transferred)); } void Connection::stop() { - _socket.close(); + std::cout << "Server::Connection::stop, Stopping connection" << std::endl; + _socket.close(); } void -Connection::handleRead(const boost::system::error_code& error, std::size_t bytes_transferred) +Connection::handleReadHeader(const boost::system::error_code& error, std::size_t bytes_transferred) { if (!error) { -/* boost::asio::async_write(_socket, reply_.to_buffers(), - boost::bind(&Connection::handleWrite, shared_from_this(), - boost::asio::placeholders::error));*/ + if (bytes_transferred != Remote::Header::size) + { + std::cerr << "bytes_transferred (" << bytes_transferred << ") != Remote::Header::size!" << std::endl; + _connectionManager.stop(shared_from_this()); + return; + } + _inputStreamBuf.commit(bytes_transferred); + + std::istream is(&_inputStreamBuf); + + Remote::Header header; + if (!header.from_istream(is)) + { + std::cerr << "Cannot read header from buffer!" << std::endl; + _connectionManager.stop(shared_from_this()); + return; + } + + std::cout << "Header received. Size = " << header.getSize() << std::endl; + + // Now read the real message + boost::asio::streambuf::mutable_buffers_type bufs = _inputStreamBuf.prepare(header.getSize()); + + boost::asio::async_read(_socket, + bufs, + boost::asio::transfer_exactly(header.getSize()), + boost::bind(&Connection::handleReadMsg, shared_from_this(), + boost::asio::placeholders::error, + boost::asio::placeholders::bytes_transferred)); - start(); } else if (error != boost::asio::error::operation_aborted) { @@ -60,21 +95,90 @@ Connection::handleRead(const boost::system::error_code& error, std::size_t bytes } } -void Connection::handleWrite(const boost::system::error_code& error) +void +Connection::handleReadMsg(const boost::system::error_code& error, std::size_t bytes_transferred) { if (!error) { - // Initiate graceful Connection closure. - boost::system::error_code ignored_ec; - _socket.shutdown(boost::asio::ip::tcp::socket::shutdown_both, ignored_ec); - } + _inputStreamBuf.commit(bytes_transferred); - if (error != boost::asio::error::operation_aborted) + std::istream is(&_inputStreamBuf); + std::ostream os(&_outputStreamBuf); + + std::vector responses; + Remote::ClientMessage request; + + if (!request.ParseFromIstream(&is)) + { + std::cerr << "Cannot parse request!" << std::endl; + _connectionManager.stop(shared_from_this()); + return; + } + + if (!_requestHandler.process(request, responses)) + { + std::cerr << "Cannot process request!" << std::endl; + _connectionManager.stop(shared_from_this()); + return; + } + + BOOST_FOREACH(const Remote::ServerMessage& response, responses) + { + boost::system::error_code ec; + + if (!response.SerializeToOstream(&os)) + { + std::cerr << "Cannot serialize to ostream!" << std::endl; + _connectionManager.stop(shared_from_this()); + return; + } + + std::array headerBuffer; + { + Remote::Header header; + header.setSize(_outputStreamBuf.size()); + header.to_buffer(headerBuffer); + } + + std::size_t n = boost::asio::write(_socket, + boost::asio::buffer(headerBuffer), + boost::asio::transfer_exactly(Remote::Header::size), + ec); + if (ec) + { + std::cerr << "cannot write header: " << error.message() << std::endl; + _connectionManager.stop(shared_from_this()); + } + + // Now send serialized payload + n = boost::asio::write(_socket, + _outputStreamBuf.data(), + boost::asio::transfer_exactly(_outputStreamBuf.size()), + ec); + + _outputStreamBuf.consume(n); + assert(n == _outputStreamBuf.size()); + + if (ec) + { + std::cerr << "cannot write msg: " << error.message() << std::endl; + _connectionManager.stop(shared_from_this()); + } + } + start(); + + // Initiate graceful Connection closure. + // boost::system::error_code ignored_ec; + // _socket.shutdown(boost::asio::ip::tcp::socket::shutdown_both, ignored_ec); + + } + else if (error != boost::asio::error::operation_aborted) { - std::cerr << "Connection::handleWrite: " << error.message() << std::endl; + std::cerr << "Connection::handleRead: " << error.message() << std::endl; _connectionManager.stop(shared_from_this()); } } + } // namespace Server } // namespace Remote diff --git a/remote/server/Connection.hpp b/remote/server/Connection.hpp index 0c99ab10..5a589178 100644 --- a/remote/server/Connection.hpp +++ b/remote/server/Connection.hpp @@ -39,11 +39,11 @@ class Connection : public std::enable_shared_from_this private: /// Handle completion of a read operation. - void handleRead(const boost::system::error_code& e, - std::size_t bytes_transferred); + void handleReadHeader(const boost::system::error_code& e, + std::size_t bytes_transferred); - /// Handle completion of a write operation. - void handleWrite(const boost::system::error_code& e); + void handleReadMsg(const boost::system::error_code& e, + std::size_t bytes_transferred); /// Socket for the connection. boost::asio::ip::tcp::socket _socket; @@ -54,18 +54,8 @@ class Connection : public std::enable_shared_from_this /// The handler used to process the incoming requests. RequestHandler& _requestHandler; - // TODO use streambuffers - - /// The incoming request. -// request request_; - - /// The parser for the incoming request. -// request_parser request_parser_; - - /// The reply to be sent back to the client. -// reply reply_; - - std::array _headerBuffer; + boost::asio::streambuf _inputStreamBuf; + boost::asio::streambuf _outputStreamBuf; }; diff --git a/remote/server/RequestHandler.cpp b/remote/server/RequestHandler.cpp index 16583e26..ff0d78b6 100644 --- a/remote/server/RequestHandler.cpp +++ b/remote/server/RequestHandler.cpp @@ -5,10 +5,39 @@ namespace Remote { namespace Server { RequestHandler::RequestHandler(boost::filesystem::path dbPath) -: _db( dbPath ) +: _db( dbPath ), +_audioCollectionRequestHandler(_db) { } +bool +RequestHandler::process(const ClientMessage& request, std::vector responses) +{ + std::cout << "TODO: process request!" << std::endl; + + switch(request.type()) + { + + case ClientMessage_Type_AuthRequest: + + break; + case ClientMessage_Type_AudioCollectionRequest: + if (request.has_audio_collection_request()) + return _audioCollectionRequestHandler.process(request.audio_collection_request(), responses); + else + std::cerr << "Malformed AudioCollectionRequest message!" << std::endl; + break; + + case ClientMessage_Type_MediaRequest: + break; + + default: + std::cerr << "Unhandled message type = " << request.type() << std::endl; + } + + return false; +} + } // namespace Server } // namespace Remote diff --git a/remote/server/RequestHandler.hpp b/remote/server/RequestHandler.hpp index 0a38f848..7538557e 100644 --- a/remote/server/RequestHandler.hpp +++ b/remote/server/RequestHandler.hpp @@ -3,9 +3,11 @@ #include +#include "messages/messages.pb.h" + #include "database/DatabaseHandler.hpp" -// #include "remote/messages/ +#include "AudioCollectionRequestHandler.hpp" namespace Remote { namespace Server { @@ -16,11 +18,16 @@ class RequestHandler RequestHandler(boost::filesystem::path dbPath); + bool process(const ClientMessage& request, std::vector responses); + + bool processAudioCollectionRequest(const AudioCollectionRequest& request, std::vector responses); private: DatabaseHandler _db; + AudioCollectionRequestHandler _audioCollectionRequestHandler; + }; } // namespace Server diff --git a/test/Makefile.am b/test/Makefile.am index 8b07f6b0..e29bc944 100644 --- a/test/Makefile.am +++ b/test/Makefile.am @@ -8,6 +8,7 @@ remote_SOURCES = \ $(srcdir)/TestDatabase.cpp \ $(top_srcdir)/remote/server/Server.cpp \ $(top_srcdir)/remote/server/RequestHandler.cpp \ + $(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp \ $(top_srcdir)/remote/server/ConnectionManager.cpp \ $(top_srcdir)/remote/server/Connection.cpp \ $(top_srcdir)/remote/messages/auth.pb.cc \ diff --git a/test/Makefile.in b/test/Makefile.in index 1bad117c..006b5e8d 100644 --- a/test/Makefile.in +++ b/test/Makefile.in @@ -120,6 +120,7 @@ database_integrity_LINK = $(CXXLD) $(database_integrity_CXXFLAGS) \ am_remote_OBJECTS = remote-RemoteClientServer.$(OBJEXT) \ remote-TestDatabase.$(OBJEXT) remote-Server.$(OBJEXT) \ remote-RequestHandler.$(OBJEXT) \ + remote-AudioCollectionRequestHandler.$(OBJEXT) \ remote-ConnectionManager.$(OBJEXT) remote-Connection.$(OBJEXT) \ remote-auth.pb.$(OBJEXT) remote-collection.pb.$(OBJEXT) \ remote-common.pb.$(OBJEXT) remote-media.pb.$(OBJEXT) \ @@ -497,6 +498,7 @@ remote_SOURCES = \ $(srcdir)/TestDatabase.cpp \ $(top_srcdir)/remote/server/Server.cpp \ $(top_srcdir)/remote/server/RequestHandler.cpp \ + $(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp \ $(top_srcdir)/remote/server/ConnectionManager.cpp \ $(top_srcdir)/remote/server/Connection.cpp \ $(top_srcdir)/remote/messages/auth.pb.cc \ @@ -622,6 +624,7 @@ distclean-compile: @AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/database_integrity-Track.Po@am__quote@ @AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/database_integrity-Video.Po@am__quote@ @AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/remote-Artist.Po@am__quote@ +@AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/remote-AudioCollectionRequestHandler.Po@am__quote@ @AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/remote-Checksum.Po@am__quote@ @AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/remote-Connection.Po@am__quote@ @AMDEP_TRUE@@am__include@ @am__quote@./$(DEPDIR)/remote-ConnectionManager.Po@am__quote@ @@ -965,6 +968,20 @@ remote-RequestHandler.obj: $(top_srcdir)/remote/server/RequestHandler.cpp @AMDEP_TRUE@@am__fastdepCXX_FALSE@ DEPDIR=$(DEPDIR) $(CXXDEPMODE) $(depcomp) @AMDEPBACKSLASH@ @am__fastdepCXX_FALSE@ $(AM_V_CXX@am__nodep@)$(CXX) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) $(CPPFLAGS) $(remote_CXXFLAGS) $(CXXFLAGS) -c -o remote-RequestHandler.obj `if test -f '$(top_srcdir)/remote/server/RequestHandler.cpp'; then $(CYGPATH_W) '$(top_srcdir)/remote/server/RequestHandler.cpp'; else $(CYGPATH_W) '$(srcdir)/$(top_srcdir)/remote/server/RequestHandler.cpp'; fi` +remote-AudioCollectionRequestHandler.o: $(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp +@am__fastdepCXX_TRUE@ $(AM_V_CXX)$(CXX) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) $(CPPFLAGS) $(remote_CXXFLAGS) $(CXXFLAGS) -MT remote-AudioCollectionRequestHandler.o -MD -MP -MF $(DEPDIR)/remote-AudioCollectionRequestHandler.Tpo -c -o remote-AudioCollectionRequestHandler.o `test -f '$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp' || echo '$(srcdir)/'`$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp +@am__fastdepCXX_TRUE@ $(AM_V_at)$(am__mv) $(DEPDIR)/remote-AudioCollectionRequestHandler.Tpo $(DEPDIR)/remote-AudioCollectionRequestHandler.Po +@AMDEP_TRUE@@am__fastdepCXX_FALSE@ $(AM_V_CXX)source='$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp' object='remote-AudioCollectionRequestHandler.o' libtool=no @AMDEPBACKSLASH@ +@AMDEP_TRUE@@am__fastdepCXX_FALSE@ DEPDIR=$(DEPDIR) $(CXXDEPMODE) $(depcomp) @AMDEPBACKSLASH@ +@am__fastdepCXX_FALSE@ $(AM_V_CXX@am__nodep@)$(CXX) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) $(CPPFLAGS) $(remote_CXXFLAGS) $(CXXFLAGS) -c -o remote-AudioCollectionRequestHandler.o `test -f '$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp' || echo '$(srcdir)/'`$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp + +remote-AudioCollectionRequestHandler.obj: $(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp +@am__fastdepCXX_TRUE@ $(AM_V_CXX)$(CXX) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) $(CPPFLAGS) $(remote_CXXFLAGS) $(CXXFLAGS) -MT remote-AudioCollectionRequestHandler.obj -MD -MP -MF $(DEPDIR)/remote-AudioCollectionRequestHandler.Tpo -c -o remote-AudioCollectionRequestHandler.obj `if test -f '$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp'; then $(CYGPATH_W) '$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp'; else $(CYGPATH_W) '$(srcdir)/$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp'; fi` +@am__fastdepCXX_TRUE@ $(AM_V_at)$(am__mv) $(DEPDIR)/remote-AudioCollectionRequestHandler.Tpo $(DEPDIR)/remote-AudioCollectionRequestHandler.Po +@AMDEP_TRUE@@am__fastdepCXX_FALSE@ $(AM_V_CXX)source='$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp' object='remote-AudioCollectionRequestHandler.obj' libtool=no @AMDEPBACKSLASH@ +@AMDEP_TRUE@@am__fastdepCXX_FALSE@ DEPDIR=$(DEPDIR) $(CXXDEPMODE) $(depcomp) @AMDEPBACKSLASH@ +@am__fastdepCXX_FALSE@ $(AM_V_CXX@am__nodep@)$(CXX) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) $(CPPFLAGS) $(remote_CXXFLAGS) $(CXXFLAGS) -c -o remote-AudioCollectionRequestHandler.obj `if test -f '$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp'; then $(CYGPATH_W) '$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp'; else $(CYGPATH_W) '$(srcdir)/$(top_srcdir)/remote/server/AudioCollectionRequestHandler.cpp'; fi` + remote-ConnectionManager.o: $(top_srcdir)/remote/server/ConnectionManager.cpp @am__fastdepCXX_TRUE@ $(AM_V_CXX)$(CXX) $(DEFS) $(DEFAULT_INCLUDES) $(INCLUDES) $(AM_CPPFLAGS) $(CPPFLAGS) $(remote_CXXFLAGS) $(CXXFLAGS) -MT remote-ConnectionManager.o -MD -MP -MF $(DEPDIR)/remote-ConnectionManager.Tpo -c -o remote-ConnectionManager.o `test -f '$(top_srcdir)/remote/server/ConnectionManager.cpp' || echo '$(srcdir)/'`$(top_srcdir)/remote/server/ConnectionManager.cpp @am__fastdepCXX_TRUE@ $(AM_V_at)$(am__mv) $(DEPDIR)/remote-ConnectionManager.Tpo $(DEPDIR)/remote-ConnectionManager.Po diff --git a/test/RemoteClientServer.cpp b/test/RemoteClientServer.cpp index 86ad3dc1..9dfb0807 100644 --- a/test/RemoteClientServer.cpp +++ b/test/RemoteClientServer.cpp @@ -44,7 +44,6 @@ class TestServer }; - class TestClient { public: @@ -56,7 +55,21 @@ class TestClient void getArtists(std::vector& artists) { + // Send request + Remote::ClientMessage msg; + msg.set_type( Remote::ClientMessage_Type_AudioCollectionRequest ); + + msg.mutable_audio_collection_request()->set_type( Remote::AudioCollectionRequest_Type_TypeGetArtistList); + msg.mutable_audio_collection_request()->mutable_get_artists()->set_preferred_batch_size(64); + + sendMsg(msg); + + // Receive responses + Remote::ServerMessage response; + recvMsg(response); + + // } @@ -66,10 +79,10 @@ class TestClient void sendMsg(const ::google::protobuf::Message& message) { + std::ostream os(&_outputStreamBuf); // Serialize message - std::string outputStr; - if (message.SerializeToString(&outputStr)) + if (message.SerializeToOstream(&os)) { // Send message header std::array headerBuffer; @@ -77,14 +90,21 @@ class TestClient // Generate header content { Remote::Header header; - header.setSize(outputStr.size()); + header.setSize(_outputStreamBuf.size()); header.to_buffer(headerBuffer); } boost::asio::write(_socket, boost::asio::buffer(headerBuffer)); // Send serialized message - boost::asio::write(_socket, boost::asio::buffer(outputStr)); + std::size_t n = boost::asio::write(_socket, + _outputStreamBuf.data(), + boost::asio::transfer_exactly(_outputStreamBuf.size())); + assert(n == _outputStreamBuf.size()); + + _outputStreamBuf.consume(n); + + std::cout << "Client: Message sent!" << std::endl; } else { @@ -96,40 +116,53 @@ class TestClient void recvMsg(::google::protobuf::Message& message) { - boost::asio::streambuf messageStreamBuffer; - std::istream responseStream(&response); + std::istream is(&_inputStreamBuf); - boost::asio::read(_socket, - messageStreamBuffer, - boost::asio::transfer_exactly(Remote::Header::size)); + { + // reserve bytes in output sequence + boost::asio::streambuf::mutable_buffers_type bufs = _inputStreamBuf.prepare(Remote::Header::size); - bool error; - Remote::Header header (Remote::Header::from_buffer(headerBuffer, error)); - if (error) + std::cout << "Client: waiting for header message!" << std::endl; + std::size_t n = boost::asio::read(_socket, + bufs, + boost::asio::transfer_exactly(Remote::Header::size)); + + assert(n == Remote::Header::size); + _inputStreamBuf.commit(n); + + std::cout << "Client: Header message received!" << std::endl; + } + + Remote::Header header; + if (!header.from_istream(is)) throw std::runtime_error("Cannot read header from buffer!"); - if (header.getSize() > _inputBuffer.size()) - throw std::runtime_error("Input buffer too small!"); - // Read in a stream buffer { - boost::asio::streambuf response; - boost::asio::read(_socket, - response, - boost::asio::transfer_exactly(header.getSize())); + // reserve bytes in output sequence + boost::asio::streambuf::mutable_buffers_type bufs = _inputStreamBuf.prepare(header.getSize()); + std::cout << "Client: waiting for message!" << std::endl; + std::size_t n = boost::asio::read(_socket, + bufs, + boost::asio::transfer_exactly(header.getSize())); - if (!message.ParseFromIstream(&responseStream)) + assert(n == header.getSize()); + _inputStreamBuf.commit(n); + + std::cout << "Client: message received!" << std::endl; + + if (!message.ParseFromIstream(&is)) throw std::runtime_error("message.ParseFromIstream failed!"); + } } boost::asio::io_service _ioService; boost::asio::ip::tcp::socket _socket; - boost::streambuf _inputStreamBuf; - - std::array _inputBuffer; + boost::asio::streambuf _inputStreamBuf; + boost::asio::streambuf _outputStreamBuf; }; diff --git a/test/TestDatabase.cpp b/test/TestDatabase.cpp index 67bb8f0f..1c1ff5d8 100644 --- a/test/TestDatabase.cpp +++ b/test/TestDatabase.cpp @@ -9,7 +9,7 @@ DatabaseHandler* create() boost::filesystem::path p ("test_db"); // Remove previous db - boost::filesystem::remove(p); +// boost::filesystem::remove(p); DatabaseHandler* db = new DatabaseHandler(p);