WIP, Remote Client/Server, fist working transcode

This commit is contained in:
emeric
2014-05-28 14:49:45 +02:00
parent 61341a090a
commit d6ad178bdb
19 changed files with 2690 additions and 858 deletions
+7 -2
View File
@@ -31,13 +31,18 @@ class Header
bool from_buffer(const std::array<unsigned char, size>& buffer)
{
if (decode32(&buffer[0]) != _magic)
if (decode32(&buffer[0]) != _magic) {
std::cerr << "Header: bad magic" << std::endl;
return false;
}
else
{
_size = decode32(&buffer[4]);
return _size < _maxSize;
if (_size > _maxSize)
std::cerr << "Header: msg too big (" << _size << ")!" << std::endl;
return _size <= _maxSize;
}
}
+1141 -425
View File
File diff suppressed because it is too large Load Diff
+845 -363
View File
File diff suppressed because it is too large Load Diff
+44 -23
View File
@@ -2,9 +2,14 @@ import "common.proto";
package Remote;
enum CodecType
enum AudioCodecType
{
CodecTypeOGG = 1;
CodecTypeOGA = 1;
}
enum VideoCodecType
{
CodecTypeOGV = 1;
}
@@ -13,24 +18,41 @@ message MediaRequest
message Prepare
{
required int64 id = 1; // Id if the media
required uint32 offset_secs = 2; // Offset in seconds
optional CodecType codec_type = 3; // Codec of the media
message Audio
{
required int64 track_id = 1; // Id if the media
optional AudioCodecType codec_type = 2;
optional uint32 bitrate = 3;
optional uint32 stream_idx = 4;
optional uint32 offset_secs = 5;
}
optional uint32 audio_bitrate = 4;
message Video
{
required int64 video_id = 1; // Id if the media
optional VideoCodecType codec_type = 2;
optional uint32 bitrate = 3;
optional uint32 offset_secs = 4;
optional uint32 audio_stream_idx = 5;
optional uint32 video_stream_idx = 6;
optional uint32 subtitle_stream_idx = 7;
}
// Video specifics
optional uint32 video_bitrate = 5;
optional uint32 audio_stream_idx = 6;
optional uint32 video_stream_idx = 7;
optional uint32 subtitle_stream_idx = 8;
enum Type {
AudioRequest = 1;
VideoRequest = 2;
}
required Type type = 1;
optional Audio audio = 2;
optional Video video = 3;
}
message GetPart
{
required uint32 requested_data_size = 1; // Amount of data requested. May receive more or less
required uint32 requested_data_size = 1; // Amount of data requested. May receive more or less. May receive 0 byte if complete
}
message Terminate
@@ -41,11 +63,11 @@ message MediaRequest
enum Type
{
TypeMediaPrepare = 1;
TypeMediaGet = 2;
TypeMediaGetPart = 2;
TypeMediaTerminate = 3;
}
required Type request_type = 1;
required Type type = 1;
optional Prepare prepare = 2;
optional GetPart get_part = 3;
@@ -55,21 +77,20 @@ message MediaRequest
message MediaResponse
{
message MediaPart
message Part
{
required uint64 byte_offset = 1;
repeated bytes data = 2;
required bytes data = 1;
}
enum Type
{
TypeMediaPart = 1;
TypeError = 1;
TypePart = 2;
}
required Error error = 1;
required Type type = 1;
required Type response_type = 2;
optional MediaPart part = 3;
optional Error error = 2;
optional Part part = 3;
}
@@ -22,47 +22,47 @@ AudioCollectionRequestHandler::process(const AudioCollectionRequest& request, Au
switch (request.type())
{
case AudioCollectionRequest_Type_TypeGetGenreList:
case AudioCollectionRequest::TypeGetGenreList:
if (request.has_get_genres())
{
res = processGetGenres(request.get_genres(), *response.mutable_genre_list());
if (res)
response.set_type(AudioCollectionResponse_Type_TypeArtistList);
response.set_type(AudioCollectionResponse::TypeArtistList);
}
else
std::cerr << "Bad AudioCollectionRequest_Type_TypeGetGenreList" << std::endl;
break;
case AudioCollectionRequest_Type_TypeGetArtistList:
case AudioCollectionRequest::TypeGetArtistList:
if (request.has_get_artists())
{
res = processGetArtists(request.get_artists(), *response.mutable_artist_list());
if (res)
response.set_type(AudioCollectionResponse_Type_TypeArtistList);
response.set_type(AudioCollectionResponse::TypeArtistList);
}
else
std::cerr << "Bad AudioCollectionRequest_Type_TypeGetArtistList" << std::endl;
break;
case AudioCollectionRequest_Type_TypeGetReleaseList:
case AudioCollectionRequest::TypeGetReleaseList:
if (request.has_get_releases())
{
res = processGetReleases(request.get_releases(), *response.mutable_release_list());
if (res)
response.set_type(AudioCollectionResponse_Type_TypeReleaseList);
response.set_type(AudioCollectionResponse::TypeReleaseList);
}
else
std::cerr << "Bad AudioCollectionRequest_Type_TypeGetReleaseList" << std::endl;
break;
case AudioCollectionRequest_Type_TypeGetTrackList:
case AudioCollectionRequest::TypeGetTrackList:
if (request.has_get_tracks())
{
res = processGetTracks(request.get_tracks(), *response.mutable_track_list());
if (res)
response.set_type(AudioCollectionResponse_Type_TypeTrackList);
response.set_type(AudioCollectionResponse::TypeTrackList);
}
else
@@ -27,7 +27,7 @@ class AudioCollectionRequestHandler
static const std::size_t _maxListArtists = 256;
static const std::size_t _maxListGenres = 256;
static const std::size_t _maxListReleases = 256;
static const std::size_t _maxListTracks = 256;
static const std::size_t _maxListTracks = 1024;
};
} // namespace Remote
+12 -4
View File
@@ -18,7 +18,8 @@ namespace Server {
Connection::Connection(boost::asio::ip::tcp::socket socket,
ConnectionManager& manager,
RequestHandler& handler)
: _socket(std::move(socket)),
: _closing(false),
_socket(std::move(socket)),
_connectionManager(manager),
_requestHandler(handler)
{
@@ -47,8 +48,15 @@ Connection::start()
void
Connection::stop()
{
std::cout << "Server::Connection::stop, Stopping connection" << std::endl;
_socket.close();
if (!_closing)
{
_closing = true;
std::cout << "Server::Connection::stop, Stopping connection " << this << std::endl;
_socket.close();
std::cout << "Server::Connection::stop, connection stopped " << this << std::endl;
}
else
std::cout << "Close in progress..." << std::endl;
}
void
@@ -88,7 +96,7 @@ Connection::handleReadHeader(const boost::system::error_code& error, std::size_t
}
else if (error != boost::asio::error::operation_aborted)
{
std::cerr << "Connection::handleRead: " << error.message() << std::endl;
std::cerr << "Connection::handleReadHeader: " << error.message() << std::endl;
_connectionManager.stop(shared_from_this());
}
}
+2
View File
@@ -38,6 +38,8 @@ class Connection : public std::enable_shared_from_this<Connection>
void stop();
private:
bool _closing;
/// Handle completion of a read operation.
void handleReadHeader(const boost::system::error_code& e,
std::size_t bytes_transferred);
+176
View File
@@ -0,0 +1,176 @@
#include "MediaRequestHandler.hpp"
#include "database/AudioTypes.hpp"
namespace Remote {
namespace Server {
MediaRequestHandler::MediaRequestHandler(DatabaseHandler& db)
: _db(db)
{}
bool
MediaRequestHandler::process(const MediaRequest& request, MediaResponse& response)
{
bool res = false;
switch (request.type())
{
case MediaRequest::TypeMediaPrepare:
if (request.has_prepare())
{
if (request.prepare().has_audio())
res = processAudioPrepare(request.prepare().audio(), response);
else if (request.prepare().has_video())
;// TODO;
else
std::cerr << "Bad MediaRequest::TypeMediaPrepare!" << std::endl;
}
else
std::cerr << "Bad MediaRequest::TypeMediaPrepare!" << std::endl;
break;
case MediaRequest::TypeMediaGetPart:
if (request.has_get_part())
res = processGetPart(request.get_part(), response);
else
std::cerr << "Bad MediaRequest::TypeMediaGet!" << std::endl;
break;
case MediaRequest::TypeMediaTerminate:
if (request.has_terminate())
res = processTerminate(request.terminate(), response);
else
std::cerr << "Bad MediaRequest::TypeMediaTerminate!" << std::endl;
break;
default:
std::cerr << "Unhandled MediaRequest type = " << request.type() << std::endl;
}
return res;
}
bool
MediaRequestHandler::processAudioPrepare(const MediaRequest::Prepare::Audio& request, MediaResponse& response)
{
// TODO, get user default values
Transcode::Format::Encoding format = Transcode::Format::OGA;
std::size_t bitrate = 128000;
if (request.has_codec_type())
{
switch( request.codec_type())
{
case AudioCodecType::CodecTypeOGA:
format = Transcode::Format::OGA;
break;
default:
std::cerr << "Unhandled codec type = " << request.codec_type() << std::endl;
return false;
}
}
if (request.has_bitrate())
bitrate = request.bitrate();
if (_transcoder)
{
response.mutable_error()->set_error(true);
response.mutable_error()->set_message("Transcode already in progress");
response.set_type(MediaResponse::TypeError);
return true;
}
Wt::Dbo::Transaction transaction( _db.getSession());
Track::pointer track = Track::getById( _db.getSession(), request.track_id() );
if (!track)
{
response.mutable_error()->set_error(true);
response.mutable_error()->set_message("Cannot find requested track!");
response.set_type(MediaResponse::TypeError);
return true;
}
try
{
Transcode::InputMediaFile inputFile(track->getPath());
Transcode::Parameters parameters(inputFile, Transcode::Format::get( format ), bitrate);
_transcoder = std::make_shared<Transcode::AvConvTranscoder>( parameters );
response.mutable_error()->set_error(false);
response.mutable_error()->set_message("");
response.set_type(MediaResponse::TypeError);
}
catch(std::exception& e)
{
std::cerr << "Caught exception: " << e.what() << std::endl;
response.mutable_error()->set_error(true);
response.mutable_error()->set_message("exception: " + std::string(e.what()));
response.set_type(MediaResponse::TypeError);
}
return true;
}
bool
MediaRequestHandler::processGetPart(const MediaRequest::GetPart& request, MediaResponse& response)
{
std::size_t dataSize = request.requested_data_size();
if (dataSize > _maxPartSize)
dataSize = _maxPartSize;
if (!_transcoder)
{
response.mutable_error()->set_error(true);
response.mutable_error()->set_message("No transcoder set!");
response.set_type(MediaResponse::TypeError);
return true;
}
while (!_transcoder->isComplete() && _transcoder->getOutputData().size() < dataSize)
_transcoder->process();
std::cout << "MediaRequestHandler::processGetPart, isComplete = " << std::boolalpha << _transcoder->isComplete() << ", size = " << _transcoder->getOutputData().size() << std::endl;
Transcode::AvConvTranscoder::data_type::iterator itEnd;
if (_transcoder->getOutputData().size() > dataSize)
itEnd = _transcoder->getOutputData().begin() + dataSize;
else
itEnd = _transcoder->getOutputData().end();
response.set_type(MediaResponse::TypePart);
std::copy(_transcoder->getOutputData().begin(), itEnd, std::back_inserter(*response.mutable_part()->mutable_data()));
// Consume sent bytes
_transcoder->getOutputData().erase(_transcoder->getOutputData().begin(), itEnd);
return true;
}
bool
MediaRequestHandler::processTerminate(const MediaRequest::Terminate& /*request*/, MediaResponse& response)
{
std::cout << "MediaRequestHandler: resetting transcoder" << std::endl;
_transcoder.reset();
assert(!_transcoder);
response.mutable_error()->set_error(false);
response.mutable_error()->set_message("");
response.set_type(MediaResponse::TypeError);
return true;
}
} // namespace Remote
} // namespace Server
+40
View File
@@ -0,0 +1,40 @@
#ifndef REMOTE_MEDIA_REQUEST_HANDLER
#define REMOTE_MEDIA_REQUEST_HANDLER
#include <memory>
#include "messages/media.pb.h"
#include "database/DatabaseHandler.hpp"
#include "transcode/AvConvTranscoder.hpp"
namespace Remote {
namespace Server {
class MediaRequestHandler
{
public:
MediaRequestHandler(DatabaseHandler& db);
bool process(const MediaRequest& request, MediaResponse& response);
private:
bool processAudioPrepare(const MediaRequest::Prepare::Audio& request, MediaResponse& response);
bool processGetPart(const MediaRequest::GetPart& request, MediaResponse& response);
bool processTerminate(const MediaRequest::Terminate& request, MediaResponse& response);
// bool processVideoPrepare(const AudioCollectionRequest::GetGenreList& request, AudioCollectionResponse::GenreList& response);
std::shared_ptr<Transcode::AvConvTranscoder> _transcoder;
DatabaseHandler& _db;
static const std::size_t _maxPartSize = 65536 - 128;
};
} // namespace Remote
} // namespace Server
#endif
+12 -2
View File
@@ -6,7 +6,8 @@ namespace Server {
RequestHandler::RequestHandler(boost::filesystem::path dbPath)
: _db( dbPath ),
_audioCollectionRequestHandler(_db)
_audioCollectionRequestHandler(_db),
_mediaRequestHandler(_db)
{
}
@@ -26,13 +27,22 @@ RequestHandler::process(const ClientMessage& request, ServerMessage& response)
{
res = _audioCollectionRequestHandler.process(request.audio_collection_request(), *response.mutable_audio_collection_response());
if (res)
response.set_type( ServerMessage_Type_AudioCollectionResponse);
response.set_type( ServerMessage::AudioCollectionResponse);
}
else
std::cerr << "Malformed AudioCollectionRequest message!" << std::endl;
break;
case ClientMessage_Type_MediaRequest:
if (request.has_media_request())
{
res = _mediaRequestHandler.process(request.media_request(), *response.mutable_media_response());
if (res)
response.set_type( ServerMessage::MediaResponse);
}
else
std::cerr << "Malformed AudioCollectionRequest message!" << std::endl;
break;
break;
default:
+3 -2
View File
@@ -8,6 +8,7 @@
#include "database/DatabaseHandler.hpp"
#include "AudioCollectionRequestHandler.hpp"
#include "MediaRequestHandler.hpp"
namespace Remote {
namespace Server {
@@ -26,8 +27,8 @@ class RequestHandler
DatabaseHandler _db;
AudioCollectionRequestHandler _audioCollectionRequestHandler;
AudioCollectionRequestHandler _audioCollectionRequestHandler;
MediaRequestHandler _mediaRequestHandler;
};
} // namespace Server