diff --git a/.clang-format b/.clang-format new file mode 100644 index 0000000..fbd096a --- /dev/null +++ b/.clang-format @@ -0,0 +1,6 @@ +--- +Language: Cpp +BasedOnStyle: WebKit +PointerAlignment: Right +ColumnLimit: 80 +--- diff --git a/ADGStreamerApp/Db/ADGStreamer.template b/ADGStreamerApp/Db/ADGStreamer.template new file mode 100755 index 0000000..8f3e5ec --- /dev/null +++ b/ADGStreamerApp/Db/ADGStreamer.template @@ -0,0 +1,11 @@ +include "ADBase.template" + +record(waveform, "$(P)$(R)PipelineBuilder_RBV") +{ + field(DESC, "Current GStreamer pipeline") + field(DTYP, "asynOctetRead") + field(INP, "@asyn($(PORT),$(ADDR=0),$(TIMEOUT=1))GST_PIPELINE_STRING") + field(FTVL, "CHAR") + field(NELM, "1024") + field(SCAN, "I/O Intr") +} diff --git a/ADGStreamerApp/Db/Makefile b/ADGStreamerApp/Db/Makefile index 8eb9727..2998f1f 100644 --- a/ADGStreamerApp/Db/Makefile +++ b/ADGStreamerApp/Db/Makefile @@ -7,6 +7,7 @@ include $(TOP)/configure/CONFIG # Create and install (or just install) into /db # databases, templates, substitutions like this #DB += xxx.db +DB += ADGStreamer.template #---------------------------------------------------- # If .db template is not named *.template add diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp new file mode 100755 index 0000000..4b87d60 --- /dev/null +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -0,0 +1,528 @@ +#include "ADGStreamer.h" + +static const char *driverName = "ADGStreamer"; + +ADGStreamer::ADGStreamer(const char *portName, int maxBuffers, size_t maxMemory, + int priority, int stackSize) + : ADDriver(portName, 1, 1, maxBuffers, maxMemory, 0, 0, ASYN_CANBLOCK, 1, + priority, stackSize) +{ + createParam(ADGSTPipelineBuilder, asynParamOctet, &GST_PipelineBuilder); + + initializeGStreamer(); + reportStatus("GStreamer initialized", ADStatusIdle); + + epicsThreadCreate("ADGSTAcquire", epicsThreadPriorityMedium, + epicsThreadGetStackSize(epicsThreadStackMedium), + (EPICSTHREADFUNC)acquisitionTaskC, this); +} + +ADGStreamer::~ADGStreamer() +{ + stopPipeline(); + delete pipelineBuilder_; +} + +asynStatus ADGStreamer::writeInt32(asynUser *pasynUser, epicsInt32 value) +{ + int function = pasynUser->reason; + + if (function == ADAcquire) { + if (value) { + startPipeline(); + } else { + stopPipeline(); + } + } + + return ADDriver::writeInt32(pasynUser, value); +} + +bool ADGStreamer::pipelineBuilder(PipelineBuilder *builder) +{ + if (!builder) + return false; + + if (pipelineBuilder_) + delete pipelineBuilder_; + + pipelineBuilder_ = builder; + + return true; +} + +void ADGStreamer::reportStatus(const char *message, ADStatus_t status) +{ + setStringParam(ADStatusMessage, message); + setIntegerParam(ADStatus, status); + callParamCallbacks(); +} + +void ADGStreamer::acquisitionTaskC(void *drvPvt) +{ + ADGStreamer *pPvt = static_cast(drvPvt); + pPvt->acquisitionTask(); +} + +void ADGStreamer::acquisitionTask() +{ + int acquire; + epicsTimeStamp startTime, endTime; + double elapsedTime, delay; + + while (true) { + epicsTimeGetCurrent(&startTime); + + lock(); + getIntegerParam(ADAcquire, &acquire); + if (acquire) { + processBusMessage(); + } + unlock(); + + epicsTimeGetCurrent(&endTime); + elapsedTime = epicsTimeDiffInSeconds(&endTime, &startTime); + delay = 0.1 - elapsedTime; + epicsThreadSleep(delay); + } +} + +bool ADGStreamer::initializeGStreamer() +{ + gst_init(nullptr, nullptr); + return true; +} + +bool ADGStreamer::startPipeline() +{ + if (pipeline_) { + return false; + } + + reportStatus("Starting pipeline", ADStatusInitializing); + + if (!createPipeline()) { + return false; + } + + firstSample_ = true; + + GstStateChangeReturn ret; + ret = gst_element_set_state(pipeline_, GST_STATE_PLAYING); + + if (ret == GST_STATE_CHANGE_FAILURE) { + reportStatus("Could not start GStreamer pipeline", ADStatusError); + return false; + } + + GstState state; + gst_element_get_state(pipeline_, &state, nullptr, GST_CLOCK_TIME_NONE); + + if (state != GST_STATE_PLAYING) { + reportStatus("Could not start GStreamer pipeline", ADStatusError); + return false; + } + + reportStatus("Pipeline started", ADStatusAcquire); + + return true; +} + +bool ADGStreamer::stopPipeline() +{ + reportStatus("Stopping pipeline", ADStatusAborting); + if (pipeline_) { + GstStateChangeReturn ret; + ret = gst_element_set_state(pipeline_, GST_STATE_NULL); + + if (ret == GST_STATE_CHANGE_FAILURE) { + reportStatus("Failed to stop pipeline", ADStatusError); + return false; + } + + GstState state; + gst_element_get_state(pipeline_, &state, nullptr, GST_CLOCK_TIME_NONE); + + if (state != GST_STATE_NULL) { + reportStatus("Failed to stop pipeline", ADStatusError); + return false; + } + } + + if (sink_) { + gst_object_unref(sink_); + sink_ = nullptr; + } + + if (bus_) { + gst_object_unref(bus_); + bus_ = nullptr; + } + + if (pipeline_) { + gst_object_unref(pipeline_); + pipeline_ = nullptr; + } + + reportStatus("Pipeline stopped", ADStatusIdle); + + return true; +} + +bool ADGStreamer::createPipeline() +{ + if (!pipelineBuilder_) { + reportStatus("Pipeline builder not configured", ADStatusError); + return false; + } + + reportStatus("Creating pipeline", ADStatusInitializing); + + std::string pipeline = pipelineBuilder_->build(); + + setStringParam(GST_PipelineBuilder, pipeline.c_str()); + callParamCallbacks(); + + GError *error = nullptr; + pipeline_ = gst_parse_launch(pipeline.c_str(), &error); + + if (!pipeline_) { + if (error) { + reportStatus(error->message, ADStatusError); + } + return false; + } + + if (error) { + reportStatus(error->message, ADStatusError); + return false; + } + + bus_ = gst_element_get_bus(pipeline_); + if (!bus_) { + reportStatus("Could not create pipeline bus", ADStatusError); + return false; + } + + sink_ = gst_bin_get_by_name(GST_BIN(pipeline_), "appsink"); + + if (!sink_) { + reportStatus("Could not find appsink", ADStatusError); + return false; + } + + gst_app_sink_set_emit_signals(GST_APP_SINK(sink_), true); + gst_app_sink_set_drop(GST_APP_SINK(sink_), true); + gst_app_sink_set_max_buffers(GST_APP_SINK(sink_), 1); + + g_signal_connect( + sink_, "new-sample", G_CALLBACK(onNewSampleCallback), this); + + reportStatus("Pipeline created", ADStatusIdle); + + return true; +} + +GstFlowReturn ADGStreamer::onNewSampleCallback( + GstAppSink *sink, gpointer userData) +{ + ADGStreamer *pPvt = static_cast(userData); + return pPvt->onNewSample(); +} + +GstFlowReturn ADGStreamer::onNewSample() +{ + GstSample *sample = nullptr; + + lock(); + + if (!sink_) { + return GST_FLOW_ERROR; + } + + sample = gst_app_sink_pull_sample(GST_APP_SINK(sink_)); + + if (!sample) { + return GST_FLOW_ERROR; + } + + processSample(sample); + gst_sample_unref(sample); + + unlock(); + + return GST_FLOW_OK; +} + +bool ADGStreamer::processSample(GstSample *sample) +{ + GstBuffer *buffer = nullptr; + GstMapInfo map; + int imageCounter_; + int numImagesCounter_; + int arrayCallbacks_; + NDArray *pImage_; + epicsTimeStamp startTime; + + epicsTimeGetCurrent(&startTime); + + buffer = gst_sample_get_buffer(sample); + + if (!buffer) { + return false; + } + + GstCaps *caps = gst_sample_get_caps(sample); + + if (caps && firstSample_) { + if (!updateSample(caps)) { + return false; + } + firstSample_ = false; + } + + if (!gst_buffer_map(buffer, &map, GST_MAP_READ)) { + return false; + } + + pImage_ = this->pNDArrayPool->alloc(ndims_, dims_, dataType_, 0, NULL); + + if (!pImage_) { + gst_buffer_unmap(buffer, &map); + return false; + } + + memcpy(pImage_->pData, map.data, pImage_->dataSize); + gst_buffer_unmap(buffer, &map); + + getIntegerParam(NDArrayCounter, &imageCounter_); + getIntegerParam(ADNumImagesCounter, &numImagesCounter_); + getIntegerParam(NDArrayCallbacks, &arrayCallbacks_); + imageCounter_++; + numImagesCounter_++; + setIntegerParam(NDArrayCounter, imageCounter_); + setIntegerParam(ADNumImagesCounter, numImagesCounter_); + pImage_->uniqueId = imageCounter_; + pImage_->timeStamp = startTime.secPastEpoch + startTime.nsec / 1.e9; + updateTimeStamp(&pImage_->epicsTS); + + this->getAttributes(pImage_->pAttributeList); + if (arrayCallbacks_) { + doCallbacksGenericPointer(pImage_, NDArrayData, 0); + } + pImage_->release(); + + return true; +} + +bool ADGStreamer::updateSample(GstCaps *caps) +{ + GstStructure *structure; + structure = gst_caps_get_structure(caps, 0); + + if (!structure) { + return false; + } + + gint width = 0; + gint height = 0; + gst_structure_get_int(structure, "width", &width); + gst_structure_get_int(structure, "height", &height); + + epicsInt32 imageWidth, imageHeight; + if (width > 0 && height > 0) { + imageWidth = width; + imageHeight = height; + setIntegerParam(ADSizeX, imageWidth); + setIntegerParam(ADSizeY, imageHeight); + setIntegerParam(NDArraySizeX, imageWidth); + setIntegerParam(NDArraySizeY, imageHeight); + } else { + return false; + } + + const gchar *format; + format = gst_structure_get_string(structure, "format"); + if (format) { + if (strcmp(format, "GRAY8") == 0) { + ndims_ = 2; + dims_[0] = imageWidth; + dims_[1] = imageHeight; + dims_[2] = 2; + dataType_ = NDUInt8; + setIntegerParam(NDDataType, dataType_); + setIntegerParam(NDColorMode, NDColorModeMono); + setIntegerParam(NDArraySize, imageWidth * imageHeight); + } else if (strcmp(format, "RGB") == 0) { + ndims_ = 3; + dims_[0] = 3; + dims_[1] = imageWidth; + dims_[2] = imageHeight; + dataType_ = NDUInt8; + setIntegerParam(NDDataType, dataType_); + setIntegerParam(NDColorMode, NDColorModeRGB1); + setIntegerParam(NDArraySize, imageWidth * imageHeight * 3); + } + } else { + return false; + } + + setIntegerParam(NDArrayCounter, 0); + setIntegerParam(ADNumImagesCounter, 0); + + callParamCallbacks(); + + return true; +} + +bool ADGStreamer::processBusMessage() +{ + if (!bus_) { + return false; + } + + GstMessage *message; + while ((message = gst_bus_pop(bus_)) != nullptr) { + switch (GST_MESSAGE_TYPE(message)) { + case GST_MESSAGE_ERROR: { + GError *error = nullptr; + gchar *debug = nullptr; + gst_message_parse_error(message, &error, &debug); + if (error) { + reportStatus(error->message, ADStatusError); + g_error_free(error); + } + if (debug) { + g_free(debug); + } + break; + } + + case GST_MESSAGE_EOS: { + reportStatus("End of stream", ADStatusIdle); + break; + } + + case GST_MESSAGE_STATE_CHANGED: { + if (GST_MESSAGE_SRC(message) == GST_OBJECT(pipeline_)) { + GstState oldState; + GstState newState; + GstState pendingState; + gst_message_parse_state_changed( + message, &oldState, &newState, &pendingState); + + if (newState == GST_STATE_READY) { + reportStatus("Pipeline ready", ADStatusIdle); + } else if (newState == GST_STATE_PAUSED) { + reportStatus("Pipeline paused", ADStatusWaiting); + } else if (newState == GST_STATE_PLAYING) { + reportStatus("Pipeline started", ADStatusAcquire); + } + } + + break; + } + + default: + break; + } + gst_message_unref(message); + } + return true; +} + +extern "C" int ADGStreamerDrive(const char *portName, int maxBuffers, + size_t maxMemory, int priority, int stackSize) +{ + new ADGStreamer(portName, maxBuffers, maxMemory, priority, stackSize); + return asynSuccess; +} + +static const iocshArg ADGStreamerDriveArg0 = { "Port name", iocshArgString }; +static const iocshArg ADGStreamerDriveArg1 = { "Max buffers", iocshArgInt }; +static const iocshArg ADGStreamerDriveArg2 = { "Max memory", iocshArgInt }; +static const iocshArg ADGStreamerDriveArg3 = { "Priority", iocshArgInt }; +static const iocshArg ADGStreamerDriveArg4 = { "Stack size", iocshArgInt }; + +static const iocshArg *const ADGStreamerDriveArgs[] + = { &ADGStreamerDriveArg0, &ADGStreamerDriveArg1, &ADGStreamerDriveArg2, + &ADGStreamerDriveArg3, &ADGStreamerDriveArg4 }; + +static const iocshFuncDef ADGStreamerDriveFuncDef + = { "ADGStreamerDrive", 5, ADGStreamerDriveArgs }; + +static void ADGStreamerDriveCallFunc(const iocshArgBuf *args) +{ + ADGStreamerDrive( + args[0].sval, args[1].ival, args[2].ival, args[3].ival, args[4].ival); +} + +extern "C" int ADGStreamerConfigureRTSP(const char *portName, + const char *username, const char *password, const char *host, int port, + const char *path, int protocol, int latency, int dropOnLatency, int codec) +{ + ADGStreamer *pDriver + = static_cast(findAsynPortDriver(portName)); + if (!pDriver) { + printf("ADGStreamerConfigureRTSP: Port \"%s\" not found.\n", portName); + return asynError; + } + + PipelineBuilderRTSP *builder = new PipelineBuilderRTSP(); + + builder->setAuthentication(username, password); + builder->setHost(host, port, path); + builder->setProtocol(static_cast(protocol)); + builder->setLatency(latency); + builder->setDropOnLatency(dropOnLatency); + builder->setCodec(static_cast(codec)); + + if (!pDriver->pipelineBuilder(builder)) { + return asynError; + } + + return asynSuccess; +} + +static const iocshArg ADGStreamerConfigureRTSPArg0 + = { "Port name", iocshArgString }; +static const iocshArg ADGStreamerConfigureRTSPArg1 + = { "Username", iocshArgString }; +static const iocshArg ADGStreamerConfigureRTSPArg2 + = { "Password", iocshArgString }; +static const iocshArg ADGStreamerConfigureRTSPArg3 = { "Host", iocshArgString }; +static const iocshArg ADGStreamerConfigureRTSPArg4 = { "Port", iocshArgInt }; +static const iocshArg ADGStreamerConfigureRTSPArg5 = { "Path", iocshArgString }; +static const iocshArg ADGStreamerConfigureRTSPArg6 + = { "Protocol", iocshArgInt }; +static const iocshArg ADGStreamerConfigureRTSPArg7 + = { "Latency (ms)", iocshArgInt }; +static const iocshArg ADGStreamerConfigureRTSPArg8 + = { "Drop on latency", iocshArgInt }; +static const iocshArg ADGStreamerConfigureRTSPArg9 = { "Codec", iocshArgInt }; + +static const iocshArg *const ADGStreamerConfigureRTSPArgs[] + = { &ADGStreamerConfigureRTSPArg0, &ADGStreamerConfigureRTSPArg1, + &ADGStreamerConfigureRTSPArg2, &ADGStreamerConfigureRTSPArg3, + &ADGStreamerConfigureRTSPArg4, &ADGStreamerConfigureRTSPArg5, + &ADGStreamerConfigureRTSPArg6, &ADGStreamerConfigureRTSPArg7, + &ADGStreamerConfigureRTSPArg8, &ADGStreamerConfigureRTSPArg9 }; + +static const iocshFuncDef ADGStreamerConfigureRTSPFuncDef + = { "ADGStreamerConfigureRTSP", 10, ADGStreamerConfigureRTSPArgs }; + +static void ADGStreamerConfigureRTSPCallFunc(const iocshArgBuf *args) +{ + ADGStreamerConfigureRTSP(args[0].sval, args[1].sval, args[2].sval, + args[3].sval, args[4].ival, args[5].sval, args[6].ival, args[7].ival, + args[8].ival, args[9].ival); +} + +static void ADGStreamerRegister() +{ + iocshRegister(&ADGStreamerDriveFuncDef, ADGStreamerDriveCallFunc); + iocshRegister( + &ADGStreamerConfigureRTSPFuncDef, ADGStreamerConfigureRTSPCallFunc); +} + +epicsExportRegistrar(ADGStreamerRegister); diff --git a/ADGStreamerApp/src/ADGStreamer.dbd b/ADGStreamerApp/src/ADGStreamer.dbd index e70223b..929c38e 100644 --- a/ADGStreamerApp/src/ADGStreamer.dbd +++ b/ADGStreamerApp/src/ADGStreamer.dbd @@ -1,6 +1 @@ -# provide definitions such as -#include "xxxRecord.dbd" -#device(xxx,CONSTANT,devXxxSoft,"SoftChannel") -#driver(myDriver) -#registrar(myRegistrar) -#variable(myVariable) +registrar(ADGStreamerRegister) diff --git a/ADGStreamerApp/src/ADGStreamer.h b/ADGStreamerApp/src/ADGStreamer.h new file mode 100755 index 0000000..94dd25e --- /dev/null +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -0,0 +1,65 @@ +#ifndef ADGSTREAMER_H +#define ADGSTREAMER_H + +#include +#include + +#include + +#include +#include + +#include "PipelineBuilder.h" +#include "PipelineBuilderRTSP.h" +#include "PipelineTypes.h" + +#define ADGSTPipelineBuilder "GST_PIPELINE_STRING" + +class ADGStreamer : public ADDriver { + +public: + ADGStreamer(const char *portName, int maxBuffers = 0, size_t maxMemory = 0, + int priority = 0, int stackSize = 0); + virtual ~ADGStreamer(); + + virtual asynStatus writeInt32(asynUser *pasynUser, epicsInt32 value); + + bool pipelineBuilder(PipelineBuilder *builder); + void reportStatus(const char *message, ADStatus_t status); + +protected: + int GST_PipelineBuilder; + +private: + static void acquisitionTaskC(void *drvPvt); + void acquisitionTask(); + + bool initializeGStreamer(); + + bool startPipeline(); + bool stopPipeline(); + bool createPipeline(); + + static GstFlowReturn onNewSampleCallback( + GstAppSink *sink, gpointer userData); + GstFlowReturn onNewSample(); + + bool processSample(GstSample *sample); + bool updateSample(GstCaps *caps); + + bool processBusMessage(); + + GstElement *pipeline_; + GstElement *sink_; + GstBus *bus_; + + bool firstSample_; + + int ndims_; + size_t dims_[3]; + NDDataType_t dataType_; + + PipelineBuilder *pipelineBuilder_; +}; + +#endif diff --git a/ADGStreamerApp/src/Makefile b/ADGStreamerApp/src/Makefile index 010e702..1d28ad2 100644 --- a/ADGStreamerApp/src/Makefile +++ b/ADGStreamerApp/src/Makefile @@ -17,9 +17,21 @@ DBD += ADGStreamer.dbd # specify all source files to be compiled and added to the library #ADGStreamer_SRCS += xxx +ADGStreamer_SRCS += ADGStreamer.cpp +ADGStreamer_SRCS += PipelineBuilder.cpp +ADGStreamer_SRCS += PipelineBuilderRTSP.cpp + +USR_CXXFLAGS += $(shell pkg-config --cflags gstreamer-1.0) +ADGStreamer_SYS_LIBS += gstapp-1.0 +ADGStreamer_SYS_LIBS += gstbase-1.0 +ADGStreamer_SYS_LIBS += gstreamer-1.0 +ADGStreamer_SYS_LIBS += gobject-2.0 +ADGStreamer_SYS_LIBS += glib-2.0 ADGStreamer_LIBS += $(EPICS_BASE_IOC_LIBS) +include $(ADCORE)/ADApp/commonLibraryMakefile + #=========================== include $(TOP)/configure/RULES diff --git a/ADGStreamerApp/src/PipelineBuilder.cpp b/ADGStreamerApp/src/PipelineBuilder.cpp new file mode 100755 index 0000000..ec5ce38 --- /dev/null +++ b/ADGStreamerApp/src/PipelineBuilder.cpp @@ -0,0 +1,92 @@ +#include "PipelineBuilder.h" + +#include + +PipelineBuilder::PipelineBuilder() + : leaky_(LEAKY_NONE) + , queueSize_(1) + , colorMode_(COLOR_RGB) +{ +} + +PipelineBuilder &PipelineBuilder::setQueue(int size, QueueLeaky leaky) +{ + queueSize_ = size; + leaky_ = leaky; + return *this; +} + +PipelineBuilder &PipelineBuilder::setVideoConvert(ColorModeType colorMode) +{ + colorMode_ = colorMode; + return *this; +} + +std::string PipelineBuilder::build() const +{ + std::stringstream pipeline; + + /*-------------------------------------------------------------- + * Queue + *-------------------------------------------------------------*/ + + pipeline << " ! queue"; + + pipeline << " max-size-buffers=" << queueSize_; + + pipeline << " leaky="; + + switch (leaky_) { + case LEAKY_NONE: + pipeline << "no"; + break; + + case LEAKY_UPSTREAM: + pipeline << "upstream"; + break; + + case LEAKY_DOWNSTREAM: + pipeline << "downstream"; + break; + + default: + pipeline << "no"; + break; + } + + /*-------------------------------------------------------------- + * Decoder + *-------------------------------------------------------------*/ + + pipeline << " ! decodebin"; + + /*-------------------------------------------------------------- + * Convert + *-------------------------------------------------------------*/ + + pipeline << " ! videoconvert"; + + switch (colorMode_) { + case COLOR_GRAY8: + pipeline << " ! video/x-raw,format=GRAY8"; + break; + + case COLOR_RGB: + pipeline << " ! video/x-raw,format=RGB"; + break; + + default: + pipeline << " ! video/x-raw,format=RGB"; + break; + } + + /*-------------------------------------------------------------- + * Sink + *-------------------------------------------------------------*/ + + pipeline << " ! appsink"; + pipeline << " name=appsink"; + pipeline << " sync=false"; + + return pipeline.str(); +} diff --git a/ADGStreamerApp/src/PipelineBuilder.h b/ADGStreamerApp/src/PipelineBuilder.h new file mode 100755 index 0000000..d0384c9 --- /dev/null +++ b/ADGStreamerApp/src/PipelineBuilder.h @@ -0,0 +1,22 @@ +#ifndef PIPELINE_BUILDER_H +#define PIPELINE_BUILDER_H + +#include "PipelineTypes.h" +#include + +class PipelineBuilder { +public: + PipelineBuilder(); + virtual ~PipelineBuilder() = default; + virtual std::string build() const; + + PipelineBuilder &setQueue(int queueSize, QueueLeaky leaky); + PipelineBuilder &setVideoConvert(ColorModeType colorMode); + +private: + int queueSize_; + QueueLeaky leaky_; + ColorModeType colorMode_; +}; + +#endif \ No newline at end of file diff --git a/ADGStreamerApp/src/PipelineBuilderRTSP.cpp b/ADGStreamerApp/src/PipelineBuilderRTSP.cpp new file mode 100755 index 0000000..61b07ac --- /dev/null +++ b/ADGStreamerApp/src/PipelineBuilderRTSP.cpp @@ -0,0 +1,138 @@ +#include "PipelineBuilderRTSP.h" + +#include + +PipelineBuilderRTSP::PipelineBuilderRTSP() + : PipelineBuilder() + , port_(554) + , protocol_(PROTOCOL_TCP) + , codec_(CodecH264) + , latency_(0) + , dropOnLatency_(false) +{ +} + +PipelineBuilderRTSP &PipelineBuilderRTSP::setAuthentication( + const std::string &username, const std::string &password) +{ + username_ = username; + password_ = password; + return *this; +} + +PipelineBuilderRTSP &PipelineBuilderRTSP::setHost( + const std::string &host, int port, const std::string &path) +{ + host_ = host; + port_ = port; + path_ = path; + return *this; +} + +PipelineBuilderRTSP &PipelineBuilderRTSP::setProtocol(ProtocolType protocol) +{ + protocol_ = protocol; + return *this; +} + +PipelineBuilderRTSP &PipelineBuilderRTSP::setLatency(int latency) +{ + latency_ = latency; + return *this; +} + +PipelineBuilderRTSP &PipelineBuilderRTSP::setDropOnLatency(bool enable) +{ + dropOnLatency_ = enable; + return *this; +} + +PipelineBuilderRTSP &PipelineBuilderRTSP::setCodec(CodecType codec) +{ + codec_ = codec; + return *this; +} + +std::string PipelineBuilderRTSP::build() const +{ + std::stringstream pipeline; + + /*-------------------------------------------------------------- + * Source + *-------------------------------------------------------------*/ + + pipeline << "rtspsrc"; + + pipeline << " location=rtsp://"; + + if (!username_.empty()) { + pipeline << username_; + + if (!password_.empty()) { + pipeline << ":" << password_; + } + + pipeline << "@"; + } + + pipeline << host_; + + if (port_ > 0) { + pipeline << ":" << port_; + } + + pipeline << path_; + + pipeline << " protocols=" << (protocol_ == PROTOCOL_TCP ? "tcp" : "udp"); + + pipeline << " latency=" << latency_; + + pipeline << " drop-on-latency=" << (dropOnLatency_ ? "true" : "false"); + + /*-------------------------------------------------------------- + * Depayloader + *-------------------------------------------------------------*/ + + switch (codec_) { + case CodecH264: + pipeline << " ! rtph264depay"; + break; + + case CodecH265: + pipeline << " ! rtph265depay"; + break; + + case CodecMJPEG: + pipeline << " ! rtpjpegdepay"; + break; + + default: + pipeline << " ! rtph264depay"; + break; + } + + /*-------------------------------------------------------------- + * Parser + *-------------------------------------------------------------*/ + + switch (codec_) { + case CodecH264: + pipeline << " ! h264parse config-interval=-1"; + break; + + case CodecH265: + pipeline << " ! h265parse"; + break; + + case CodecMJPEG: + break; + + default: + pipeline << " ! h264parse config-interval=-1"; + break; + } + + pipeline << PipelineBuilder::build(); + + return pipeline.str(); +} diff --git a/ADGStreamerApp/src/PipelineBuilderRTSP.h b/ADGStreamerApp/src/PipelineBuilderRTSP.h new file mode 100755 index 0000000..f3aea4f --- /dev/null +++ b/ADGStreamerApp/src/PipelineBuilderRTSP.h @@ -0,0 +1,37 @@ +#ifndef RTSP_PIPELINE_BUILDER_H +#define RTSP_PIPELINE_BUILDER_H + +#include + +#include "PipelineBuilder.h" +#include "PipelineTypes.h" + +class PipelineBuilderRTSP : public PipelineBuilder { +public: + PipelineBuilderRTSP(); + + PipelineBuilderRTSP &setAuthentication( + const std::string &username, const std::string &password); + PipelineBuilderRTSP &setHost( + const std::string &host, int port, const std::string &path); + PipelineBuilderRTSP &setProtocol(ProtocolType protocol); + PipelineBuilderRTSP &setLatency(int latency); + PipelineBuilderRTSP &setDropOnLatency(bool enable); + PipelineBuilderRTSP &setCodec(CodecType codec); + + std::string build() const override; + +private: + std::string username_; + std::string password_; + std::string host_; + std::string path_; + int port_; + + ProtocolType protocol_; + CodecType codec_; + int latency_; + bool dropOnLatency_; +}; + +#endif \ No newline at end of file diff --git a/ADGStreamerApp/src/PipelineTypes.h b/ADGStreamerApp/src/PipelineTypes.h new file mode 100755 index 0000000..c10bbe4 --- /dev/null +++ b/ADGStreamerApp/src/PipelineTypes.h @@ -0,0 +1,15 @@ +#ifndef PIPELINE_TYPES_H +#define PIPELINE_TYPES_H + +enum ProtocolType { PROTOCOL_TCP = 0, PROTOCOL_UDP = 1 }; + +enum CodecType { CodecH264 = 0, CodecH265 = 1, CodecMJPEG = 2 }; + +enum QueueLeaky { LEAKY_NONE = 0, LEAKY_UPSTREAM = 1, LEAKY_DOWNSTREAM = 2 }; + +enum ColorModeType { + COLOR_GRAY8 = 0, + COLOR_RGB = 1, +}; + +#endif diff --git a/README.md b/README.md new file mode 100755 index 0000000..fd469a4 --- /dev/null +++ b/README.md @@ -0,0 +1,14 @@ +ADGStreamer +======= + +An [EPICS][epics] [areaDetector][] driver for the [GStreamer][profilers] to acquire images from network video streams and publish them as **NDArrays** through **EPICS**. + +Currently, ADGStreamer supports **RTSP** communication only and provides configuration options for several pipeline parameters, including: protocol, latency (ms), drop frames on latency, codec, queue size (buffers) and queue leaky mode. + +[profilers]: https://gstreamer.freedesktop.org/ +[epics]: https://docs.epics-controls.org/en/latest/ +[areaDetector]: https://github.com/areaDetector/areaDetector/blob/master/README.md + +Additional information: +- [Documentation](docs/ADGStreamer/ADGStreamer.rst) +- [Release notes](RELEASE.md) \ No newline at end of file diff --git a/configure/CONFIG_SITE b/configure/CONFIG_SITE index 212485e..c6717ad 100644 --- a/configure/CONFIG_SITE +++ b/configure/CONFIG_SITE @@ -41,3 +41,11 @@ CHECK_RELEASE = YES -include $(TOP)/../CONFIG_SITE.local -include $(TOP)/configure/CONFIG_SITE.local +# Get settings from AREA_DETECTOR, so we only have to configure once for all detectors if we want to +-include $(AREA_DETECTOR)/configure/CONFIG_SITE +-include $(AREA_DETECTOR)/configure/CONFIG_SITE.$(EPICS_HOST_ARCH) +-include $(AREA_DETECTOR)/configure/CONFIG_SITE.$(EPICS_HOST_ARCH).Common +ifdef T_A + -include $(AREA_DETECTOR)/configure/CONFIG_SITE.Common.$(T_A) + -include $(AREA_DETECTOR)/configure/CONFIG_SITE.$(EPICS_HOST_ARCH).$(T_A) +endif diff --git a/configure/RELEASE b/configure/RELEASE index ac5fcdd..ed39d96 100644 --- a/configure/RELEASE +++ b/configure/RELEASE @@ -39,4 +39,6 @@ EPICS_BASE = /opt/epics/base # without having to modify this file directly. -include $(TOP)/../RELEASE.local -include $(TOP)/../RELEASE.$(EPICS_HOST_ARCH).local +-include $(TOP)/../configure/RELEASE_LIBS_INCLUDE +-include $(TOP)/RELEASE.local -include $(TOP)/configure/RELEASE.local diff --git a/docs/ADGStreamer/ADGStreamer.rst b/docs/ADGStreamer/ADGStreamer.rst new file mode 100644 index 0000000..39c6f17 --- /dev/null +++ b/docs/ADGStreamer/ADGStreamer.rst @@ -0,0 +1,151 @@ +ADGStreamer +=========== + +Overview +-------- + +ADGStreamer is an AreaDetector driver that uses the GStreamer multimedia +framework to acquire images from IP cameras and other streaming sources. + +Unlike most existing AreaDetector drivers, ADGStreamer separates the +communication protocol from the driver implementation through the +PipelineBuilder architecture. This allows different streaming protocols +to share the same acquisition logic while implementing only the +protocol-specific pipeline stages. + +The driver automatically detects image width, height, color mode and +data type from the received stream and publishes the images as NDArrays. + +Features +-------- + +The current implementation provides the following features: + +- Native AreaDetector driver. +- GStreamer-based acquisition. +- Automatic image size detection. +- Automatic color mode detection. +- Automatic NDDataType detection. +- PipelineBuilder architecture. +- RTSP support. + +Future versions are expected to support additional communication +protocols including HTTP, RTMP and SRT. + +Architecture +------------ + +The driver is divided into two independent layers. + +The first layer is responsible for the communication protocol. + +The second layer contains the common image processing stages shared by +all protocols. + +:: + + +-----------------------------+ + | ADDriver | + +-------------+---------------+ + | + v + +---------------+ + | ADGStreamer | + +-------+-------+ + | + v + +------------------+ + | PipelineBuilder | + +--------+---------+ + | + +-----------+-----------+ + | | + v v + PipelineBuilderRTSP PipelineBuilderHTTP (not yet supported) + | | + +-----------+-----------+ + | + v + GStreamer Pipeline + | + v + AppSink + | + v + NDArrayPool + +Supported Protocols +------------------- + +================== =========== +Protocol Status +================== =========== +RTSP Supported +HTTP Planned +RTMP Planned +SRT Planned +================== =========== + +Requirements +------------ + +The driver requires GStreamer 1.x together with the standard plugin +packages. + +Typical Debian installation: + +:: + + sudo apt install \ + libgstreamer1.0-dev \ + libgstreamer-plugins-base1.0-dev \ + gstreamer1.0-tools \ + gstreamer1.0-plugins-base \ + gstreamer1.0-plugins-good \ + gstreamer1.0-plugins-bad \ + gstreamer1.0-plugins-ugly \ + gstreamer1.0-libav + +Driver Parameters +----------------- + +The driver uses the standard AreaDetector parameters. + +Additional parameters include: + +======================= =============================================== +Parameter Description +======================= =============================================== +PipelineBuilder_RBV Generated GStreamer pipeline +======================= =============================================== + +IOC Configuration +----------------- + +The driver is configured using IOC shell commands. + +Example: + +:: + + ADGStreamerDrive( + Port name, + Max buffers, + Max memory, + Thread priority, + Thread stack size, + ) + + ADGStreamerConfigureRTSP( + Port name, + Username, + Password, + Host, + Port, + Path, + Protocol, + latency, + Drop frames on latency, + Codec, + Queue size, + Queue leaky mode) diff --git a/iocs/ADGStreamerIOC/ADGStreamerApp/src/ADGStreamerMain.cpp b/iocs/ADGStreamerIOC/ADGStreamerApp/src/ADGStreamerMain.cpp index e3e5776..3e9af4f 100644 --- a/iocs/ADGStreamerIOC/ADGStreamerApp/src/ADGStreamerMain.cpp +++ b/iocs/ADGStreamerIOC/ADGStreamerApp/src/ADGStreamerMain.cpp @@ -2,22 +2,21 @@ /* Author: Marty Kraimer Date: 17MAR2000 */ #include +#include #include -#include #include -#include #include "epicsExit.h" #include "epicsThread.h" #include "iocsh.h" -int main(int argc,char *argv[]) +int main(int argc, char *argv[]) { - if(argc>=2) { + if (argc >= 2) { iocsh(argv[1]); epicsThreadSleep(.2); } iocsh(NULL); epicsExit(0); - return(0); + return (0); } diff --git a/iocs/ADGStreamerIOC/ADGStreamerApp/src/Makefile b/iocs/ADGStreamerIOC/ADGStreamerApp/src/Makefile index 4e64605..0315f96 100644 --- a/iocs/ADGStreamerIOC/ADGStreamerApp/src/Makefile +++ b/iocs/ADGStreamerIOC/ADGStreamerApp/src/Makefile @@ -9,6 +9,7 @@ include $(TOP)/configure/CONFIG # Build the IOC application PROD_IOC = ADGStreamer +PROD_NAME = $(PROD_IOC) # ADGStreamer.dbd will be created and installed DBD += ADGStreamer.dbd @@ -17,9 +18,11 @@ ADGStreamer_DBD += base.dbd # Include dbd files from all support applications: #ADGStreamer_DBD += xxx.dbd +ADGStreamer_DBD += ADGStreamer.dbd # Add all the support libraries needed by this IOC #ADGStreamer_LIBS += xxx +ADGStreamer_LIBS += ADGStreamer # ADGStreamer_registerRecordDeviceDriver.cpp derives from ADGStreamer.dbd ADGStreamer_SRCS += ADGStreamer_registerRecordDeviceDriver.cpp @@ -34,6 +37,8 @@ ADGStreamer_SRCS_vxWorks += -nil- # Finally link to the EPICS Base libraries ADGStreamer_LIBS += $(EPICS_BASE_IOC_LIBS) +include $(ADCORE)/ADApp/commonDriverMakefile + #=========================== include $(TOP)/configure/RULES diff --git a/iocs/ADGStreamerIOC/configure/CONFIG_SITE b/iocs/ADGStreamerIOC/configure/CONFIG_SITE index 212485e..c6717ad 100644 --- a/iocs/ADGStreamerIOC/configure/CONFIG_SITE +++ b/iocs/ADGStreamerIOC/configure/CONFIG_SITE @@ -41,3 +41,11 @@ CHECK_RELEASE = YES -include $(TOP)/../CONFIG_SITE.local -include $(TOP)/configure/CONFIG_SITE.local +# Get settings from AREA_DETECTOR, so we only have to configure once for all detectors if we want to +-include $(AREA_DETECTOR)/configure/CONFIG_SITE +-include $(AREA_DETECTOR)/configure/CONFIG_SITE.$(EPICS_HOST_ARCH) +-include $(AREA_DETECTOR)/configure/CONFIG_SITE.$(EPICS_HOST_ARCH).Common +ifdef T_A + -include $(AREA_DETECTOR)/configure/CONFIG_SITE.Common.$(T_A) + -include $(AREA_DETECTOR)/configure/CONFIG_SITE.$(EPICS_HOST_ARCH).$(T_A) +endif diff --git a/iocs/ADGStreamerIOC/configure/RELEASE b/iocs/ADGStreamerIOC/configure/RELEASE index ac5fcdd..4b851c4 100644 --- a/iocs/ADGStreamerIOC/configure/RELEASE +++ b/iocs/ADGStreamerIOC/configure/RELEASE @@ -20,6 +20,7 @@ # Variables may be used before their values have been set. # Build variables that are NOT used in paths should be set in # the CONFIG_SITE file. +ADGSTREAMER = $(TOP)/../.. # Variables and paths to dependent modules: #MODULES = /path/to/modules @@ -39,4 +40,6 @@ EPICS_BASE = /opt/epics/base # without having to modify this file directly. -include $(TOP)/../RELEASE.local -include $(TOP)/../RELEASE.$(EPICS_HOST_ARCH).local +-include $(TOP)/../../../configure/RELEASE_PRODS_INCLUDE +-include $(TOP)/RELEASE.local -include $(TOP)/configure/RELEASE.local diff --git a/iocs/ADGStreamerIOC/iocBoot/iocADGStreamer/st.cmd b/iocs/ADGStreamerIOC/iocBoot/iocADGStreamer/st.cmd old mode 100644 new mode 100755 index 3510639..5567209 --- a/iocs/ADGStreamerIOC/iocBoot/iocADGStreamer/st.cmd +++ b/iocs/ADGStreamerIOC/iocBoot/iocADGStreamer/st.cmd @@ -1,18 +1,73 @@ #!../../bin/linux-x86_64/ADGStreamer -#- You may have to change ADGStreamer to something else -#- everywhere it appears in this file +< envPaths -#< envPaths +# Prefix +epicsEnvSet("P", "SWC:A:ADGST01:") +# The port name for the detector +epicsEnvSet("PORT", "ADGStreamer") +# The queue size for all plugins +epicsEnvSet("QSIZE","400") +# The search path for database files +epicsEnvSet("EPICS_DB_INCLUDE_PATH", "${ADCORE}/db:${ADGSTREAMER}/db") ## Register all support components -dbLoadDatabase "../../dbd/ADGStreamer.dbd" -ADGStreamer_registerRecordDeviceDriver(pdbbase) +dbLoadDatabase("../../dbd/ADGStreamer.dbd") +ADGStreamer_registerRecordDeviceDriver(pdbbase) -## Load record instances -#dbLoadRecords("../../db/ADGStreamer.db","user=root") +# ADGStreamerDrive +# Port name +# Max buffers +# Max memory +# Thread priority +# Thread stack size +ADGStreamerDrive("$(PORT)", 0, 0, 0, 0) +dbLoadRecords("ADGStreamer.template", "P=${P},R=cam1:,PORT=${PORT},ADDR=0,TIMEOUT=1") + +# ADGStreamerConfigureRTSP +# Port name +# Username +# Password +# Host +# Port +# Path +# Protocol: 0 = TCP, 1 = UDP +# RTSP latency (ms) +# Drop frames on latency: 0 = false, 1 = true +# Codec: 0 = H.264, 1 = H.265, 2 = MJPEG +# Queue size (buffers) +# Queue leaky mode: 0 = none, 1 = upstream, 2 = downstream +ADGStreamerConfigureRTSP("$(PORT)", "admin", "r00tr00t", "10.20.21.40", 554, "/cam/realmonitor?channel=1&subtype=0", 1, 0, 1, 0, 1, 2) + +# Create PV Access conversion plugin +NDPvaConfigure("PVA1", ${QSIZE}, 0, "${PORT}", 0, ${P}Pva1:Image, 0, 0, 0) +dbLoadRecords("NDPva.template", "P=${P}, R=Pva1:, PORT=PVA1, ADDR=0, TIMEOUT=1, NDARRAY_PORT=${PORT}") + +# Create an Image plugin +NDStdArraysConfigure("Image1", ${QSIZE}, 0, "${PORT}", 0, 0) +dbLoadRecords("NDStdArrays.template", "P=${P},R=image1:,PORT=Image1,NDARRAY_PORT=${PORT},ADDR=0,TIMEOUT=1,TYPE=Int8,FTVL=UCHAR,NELEMENTS=10000000") + +# Create an HDF5 file saving plugin +#NDFileHDF5Configure("FileHDF1", ${QSIZE}, 0, "${PORT}", 0) +#dbLoadRecords("NDFileHDF5.template","P=${P},R=HDF1:,PORT=FileHDF1,ADDR=0,TIMEOUT=1,NDARRAY_PORT=${PORT}") + +# Turn on asyn trace +#asynSetTraceMask("${PORT}",0,0x21) +#asynSetTraceIOMask("${PORT}",0,1) iocInit() -## Start any sequence programs -#seq sncADGStreamer,"user=root" +#Enable Array Callbacks, set Attributes file +dbpf("${P}cam1:ArrayCallbacks","Enable") + +#Enable Array Callbacks image plugin +dbpf("${P}image1:EnableCallbacks","Enable") + +#Enable Array Callbacks Pva plugin +dbpf("${P}Pva1:EnableCallbacks","Enable") + +#Enable Array Callbacks HDF1 plugin +#dbpf("${P}HDF1:EnableCallbacks","Enable") + +# Start acquisition +#dbpf("${P}cam1:Acquire","1")