Experimental support for raw TCP/IP with zero copy feature
This commit is contained in:
@@ -48,7 +48,7 @@ void ZMQImagePuller::Disconnect() {
|
||||
void ZMQImagePuller::PullerThread() {
|
||||
while (true) {
|
||||
ImagePullerOutput ret;
|
||||
ret.msg = std::make_shared<ZMQMessage>();
|
||||
ret.zmq_msg = std::make_shared<ZMQMessage>();
|
||||
bool received = false;
|
||||
while (!received) {
|
||||
if (disconnect) {
|
||||
@@ -56,7 +56,7 @@ void ZMQImagePuller::PullerThread() {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
received = socket.Receive(*ret.msg, false);
|
||||
received = socket.Receive(*ret.zmq_msg, false);
|
||||
if (!received)
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(1));
|
||||
} catch (const JFJochException &e) {
|
||||
@@ -70,9 +70,9 @@ void ZMQImagePuller::PullerThread() {
|
||||
void ZMQImagePuller::CBORThread() {
|
||||
auto ret = cbor_fifo.GetBlocking();
|
||||
|
||||
while (ret.msg) {
|
||||
while (ret.zmq_msg) {
|
||||
try {
|
||||
ret.cbor = CBORStream2Deserialize(ret.msg->data(), ret.msg->size());
|
||||
ret.cbor = CBORStream2Deserialize(ret.zmq_msg->data(), ret.zmq_msg->size());
|
||||
if (ret.cbor->msg_type == CBORImageType::END)
|
||||
logger.Info("Received END");
|
||||
|
||||
@@ -98,7 +98,7 @@ void ZMQImagePuller::RepubThread() {
|
||||
auto ret = repub_fifo.GetBlocking();
|
||||
bool repub_active = false;
|
||||
|
||||
while (ret.msg) {
|
||||
while (ret.zmq_msg) {
|
||||
try {
|
||||
if (ret.cbor->msg_type == CBORImageType::START) {
|
||||
// Start message needs to be cleaned when running republish
|
||||
@@ -112,7 +112,7 @@ void ZMQImagePuller::RepubThread() {
|
||||
logger.Info("Republish active");
|
||||
} else {
|
||||
if (repub_active)
|
||||
repub_socket->Send(ret.msg->data(), ret.msg->size(), true);
|
||||
repub_socket->Send(ret.zmq_msg->data(), ret.zmq_msg->size(), true);
|
||||
}
|
||||
} catch (const JFJochException &e) {
|
||||
logger.ErrorException(e);
|
||||
|
||||
Reference in New Issue
Block a user