#include "nng_source.h" #include "dsp/block.h" #include "dsp/block_helpers.h" #include #include #include #include #include namespace satdump { namespace ndsp { template NNGSourceBlock::NNGSourceBlock() : Block("nng_source_" + getShortTypeName(), {}, // {{"out", getTypeSampleType()}}) { } template NNGSourceBlock::~NNGSourceBlock() { } template void NNGSourceBlock::start() { logger->info("Opening TCP socket on " + std::string("tcp://" + address + ":" + std::to_string(port))); nng_sub0_open_raw(&n_sock); nng_dialer_create(&n_dialer, n_sock, std::string("tcp://" + address + ":" + std::to_string(port)).c_str()); nng_dialer_start(n_dialer, (int)0); Block::start(); } template void NNGSourceBlock::stop(bool stop_now, bool force) { Block::stop(stop_now, force); nng_dialer_close(n_dialer); nng_close(n_sock); } template bool NNGSourceBlock::work() { if (work_should_exit) { outputs[0].fifo->wait_enqueue(outputs[0].fifo->newBufferTerminator()); return true; } try { void *data = NULL; size_t lpkt_size; nng_recv(n_sock, &data, &lpkt_size, NNG_FLAG_ALLOC | NNG_FLAG_NONBLOCK); auto oblk = outputs[0].fifo->newBufferSamples(lpkt_size / sizeof(T), sizeof(T)); if (data) { oblk.size = lpkt_size / sizeof(T); memcpy(oblk.template getSamples(), data, oblk.size * sizeof(T)); nng_free(data, lpkt_size); } if (data) outputs[0].fifo->wait_enqueue(oblk); else outputs[0].fifo->free(oblk); return false; } catch (std::exception &e) { logger->error("%s", e.what()); outputs[0].fifo->wait_enqueue(outputs[0].fifo->newBufferTerminator()); return true; } } template class NNGSourceBlock; template class NNGSourceBlock; template class NNGSourceBlock; template class NNGSourceBlock; template class NNGSourceBlock; } // namespace ndsp } // namespace satdump