From 08dea92e01644dc61ca9c49c34f5c60af1b12e82 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 19 Aug 2026 17:25:03 -0300 Subject: [PATCH 01/14] Add clang format file --- .clang-format | 6 ++++++ .../ADGStreamerApp/src/ADGStreamerMain.cpp | 9 ++++----- 2 files changed, 10 insertions(+), 5 deletions(-) create mode 100644 .clang-format 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/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); } From 5f8d88dc0e0e08c5e493ef457cf514b2faf0b191 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 10:12:39 -0300 Subject: [PATCH 02/14] configure: use areaDetector configure --- configure/CONFIG_SITE | 8 ++++++++ configure/RELEASE | 2 ++ 2 files changed, 10 insertions(+) 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 From 8b55858fcfae2af828f38fa40eae427131ae6c66 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 10:24:30 -0300 Subject: [PATCH 03/14] ADGStreamerApp: add gstreamer constructor and iocsh function --- ADGStreamerApp/src/ADGStreamer.cpp | 52 ++++++++++++++++++++++++++++++ ADGStreamerApp/src/ADGStreamer.dbd | 7 +--- ADGStreamerApp/src/ADGStreamer.h | 24 ++++++++++++++ ADGStreamerApp/src/Makefile | 10 ++++++ 4 files changed, 87 insertions(+), 6 deletions(-) create mode 100755 ADGStreamerApp/src/ADGStreamer.cpp create mode 100755 ADGStreamerApp/src/ADGStreamer.h diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp new file mode 100755 index 0000000..6243235 --- /dev/null +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -0,0 +1,52 @@ +#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) +{ + initializeGStreamer(); +} + +ADGStreamer::~ADGStreamer() { } + +bool ADGStreamer::initializeGStreamer() +{ + gst_init(nullptr, nullptr); + 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); +} + +static void ADGStreamerRegister() +{ + iocshRegister(&ADGStreamerDriveFuncDef, ADGStreamerDriveCallFunc); +} + +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..5a5c4d0 --- /dev/null +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -0,0 +1,24 @@ +#ifndef ADGSTREAMER_H +#define ADGSTREAMER_H + +#include +#include + +#include + +#include +#include + +class ADGStreamer : public ADDriver { + +public: + ADGStreamer(const char *portName, int maxBuffers = 0, size_t maxMemory = 0, + int priority = 0, int stackSize = 0); + virtual ~ADGStreamer(); + +protected: +private: + bool initializeGStreamer(); +}; + +#endif diff --git a/ADGStreamerApp/src/Makefile b/ADGStreamerApp/src/Makefile index 010e702..d3af65f 100644 --- a/ADGStreamerApp/src/Makefile +++ b/ADGStreamerApp/src/Makefile @@ -17,9 +17,19 @@ DBD += ADGStreamer.dbd # specify all source files to be compiled and added to the library #ADGStreamer_SRCS += xxx +ADGStreamer_SRCS += ADGStreamer.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 From 8fe65115f3a5e217f59e3e13fd005af002d4b11b Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 16:03:21 -0300 Subject: [PATCH 04/14] ADGStreamerApp: add basic acquisition support Implemented the pipeline start and stop logic through the `startPipeline()` and `stopPipeline()` routines. When an acquisition starts, samples are received by the `onNewSampleCallback()` callback function. At this stage, the pipeline string configuration is hardcoded. --- ADGStreamerApp/src/ADGStreamer.cpp | 146 ++++++++++++++++++++++++++++- ADGStreamerApp/src/ADGStreamer.h | 13 +++ 2 files changed, 158 insertions(+), 1 deletion(-) diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp index 6243235..4a79def 100755 --- a/ADGStreamerApp/src/ADGStreamer.cpp +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -10,7 +10,22 @@ ADGStreamer::ADGStreamer(const char *portName, int maxBuffers, size_t maxMemory, initializeGStreamer(); } -ADGStreamer::~ADGStreamer() { } +ADGStreamer::~ADGStreamer() { stopPipeline(); } + +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::initializeGStreamer() { @@ -18,6 +33,135 @@ bool ADGStreamer::initializeGStreamer() return true; } +bool ADGStreamer::startPipeline() +{ + if (pipeline_) { + return false; + } + + if (!createPipeline()) { + return false; + } + + GstStateChangeReturn ret; + ret = gst_element_set_state(pipeline_, GST_STATE_PLAYING); + + if (ret == GST_STATE_CHANGE_FAILURE) { + return false; + } + + GstState state; + gst_element_get_state(pipeline_, &state, nullptr, GST_CLOCK_TIME_NONE); + + if (state != GST_STATE_PLAYING) { + return false; + } + + return true; +} + +bool ADGStreamer::stopPipeline() +{ + if (pipeline_) { + GstStateChangeReturn ret; + ret = gst_element_set_state(pipeline_, GST_STATE_NULL); + + if (ret == GST_STATE_CHANGE_FAILURE) { + return false; + } + + GstState state; + gst_element_get_state(pipeline_, &state, nullptr, GST_CLOCK_TIME_NONE); + + if (state != GST_STATE_NULL) { + return false; + } + } + + if (sink_) { + gst_object_unref(sink_); + sink_ = nullptr; + } + + if (pipeline_) { + gst_object_unref(pipeline_); + pipeline_ = nullptr; + } + + return true; +} + +bool ADGStreamer::createPipeline() +{ + std::string pipeline + = std::string("rtspsrc " + "location=rtsp://admin:r00tr00t@10.20.21.40:554/cam/" + "realmonitor?channel=1&subtype=0 protocols=udp latency=0 " + "drop-on-latency=true") + + std::string(" ! rtph264depay ! h264parse config-interval=-1") + + std::string( + " ! queue max-size-buffers=1 leaky=downstream silent=true") + + std::string(" ! avdec_h264 ! videoconvert ! video/x-raw,format=RGB ! " + "appsink name=appsink sync=false"); + + GError *error = nullptr; + pipeline_ = gst_parse_launch(pipeline.c_str(), &error); + + if (!pipeline_) { + if (error) { } + return false; + } + + if (error) { + return false; + } + + sink_ = gst_bin_get_by_name(GST_BIN(pipeline_), "appsink"); + + if (!sink_) { + 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); + + 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; + } + + gst_sample_unref(sample); + + unlock(); + + return GST_FLOW_OK; +} + extern "C" int ADGStreamerDrive(const char *portName, int maxBuffers, size_t maxMemory, int priority, int stackSize) { diff --git a/ADGStreamerApp/src/ADGStreamer.h b/ADGStreamerApp/src/ADGStreamer.h index 5a5c4d0..0f622a4 100755 --- a/ADGStreamerApp/src/ADGStreamer.h +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -16,9 +16,22 @@ class ADGStreamer : public ADDriver { int priority = 0, int stackSize = 0); virtual ~ADGStreamer(); + virtual asynStatus writeInt32(asynUser *pasynUser, epicsInt32 value); + protected: private: bool initializeGStreamer(); + + bool startPipeline(); + bool stopPipeline(); + bool createPipeline(); + + static GstFlowReturn onNewSampleCallback( + GstAppSink *sink, gpointer userData); + GstFlowReturn onNewSample(); + + GstElement *pipeline_; + GstElement *sink_; }; #endif From f8b1becf49445cee2a0875dc361fbfe009afeea9 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 16:55:19 -0300 Subject: [PATCH 05/14] ADGStreamerApp: add sample processing Implemented the `processSample()` routine to process samples received by the GStreamer callback. The function extracts the image data from the `GstSample` and prepares it for delivery to the AreaDetector driver and plugins. --- ADGStreamerApp/src/ADGStreamer.cpp | 124 +++++++++++++++++++++++++++++ ADGStreamerApp/src/ADGStreamer.h | 9 +++ 2 files changed, 133 insertions(+) diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp index 4a79def..9836d1c 100755 --- a/ADGStreamerApp/src/ADGStreamer.cpp +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -43,6 +43,8 @@ bool ADGStreamer::startPipeline() return false; } + firstSample_ = true; + GstStateChangeReturn ret; ret = gst_element_set_state(pipeline_, GST_STATE_PLAYING); @@ -155,6 +157,7 @@ GstFlowReturn ADGStreamer::onNewSample() return GST_FLOW_ERROR; } + processSample(sample); gst_sample_unref(sample); unlock(); @@ -162,6 +165,127 @@ GstFlowReturn ADGStreamer::onNewSample() 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; +} + extern "C" int ADGStreamerDrive(const char *portName, int maxBuffers, size_t maxMemory, int priority, int stackSize) { diff --git a/ADGStreamerApp/src/ADGStreamer.h b/ADGStreamerApp/src/ADGStreamer.h index 0f622a4..8b1dbe0 100755 --- a/ADGStreamerApp/src/ADGStreamer.h +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -30,8 +30,17 @@ class ADGStreamer : public ADDriver { GstAppSink *sink, gpointer userData); GstFlowReturn onNewSample(); + bool processSample(GstSample *sample); + bool updateSample(GstCaps *caps); + GstElement *pipeline_; GstElement *sink_; + + bool firstSample_; + + int ndims_; + size_t dims_[3]; + NDDataType_t dataType_; }; #endif From 15fdbffabe43ec8e2478fbfff3ec79aecb813227 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 13:41:07 -0300 Subject: [PATCH 06/14] ADGStreamerApp: add epicsThread An epicsThread has been added to support the future implementation of a GstBus monitoring. --- ADGStreamerApp/src/ADGStreamer.cpp | 39 ++++++++++++++++++++++++++++++ ADGStreamerApp/src/ADGStreamer.h | 3 +++ 2 files changed, 42 insertions(+) diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp index 9836d1c..518cd96 100755 --- a/ADGStreamerApp/src/ADGStreamer.cpp +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -8,6 +8,12 @@ ADGStreamer::ADGStreamer(const char *portName, int maxBuffers, size_t maxMemory, priority, stackSize) { initializeGStreamer(); + + epicsThreadCreate("ADGSTAcquire", + epicsThreadPriorityMedium, + epicsThreadGetStackSize(epicsThreadStackMedium), + (EPICSTHREADFUNC)acquisitionTaskC, + this); } ADGStreamer::~ADGStreamer() { stopPipeline(); } @@ -27,6 +33,39 @@ asynStatus ADGStreamer::writeInt32(asynUser *pasynUser, epicsInt32 value) return ADDriver::writeInt32(pasynUser, value); } +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) + { + } + else + { + } + unlock(); + + epicsTimeGetCurrent(&endTime); + elapsedTime = epicsTimeDiffInSeconds(&endTime, &startTime); + delay = 0.1 - elapsedTime; + epicsThreadSleep(delay); + } +} + bool ADGStreamer::initializeGStreamer() { gst_init(nullptr, nullptr); diff --git a/ADGStreamerApp/src/ADGStreamer.h b/ADGStreamerApp/src/ADGStreamer.h index 8b1dbe0..291a64a 100755 --- a/ADGStreamerApp/src/ADGStreamer.h +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -20,6 +20,9 @@ class ADGStreamer : public ADDriver { protected: private: + static void acquisitionTaskC(void *drvPvt); + void acquisitionTask(); + bool initializeGStreamer(); bool startPipeline(); From b519d2d95f144a4f74cfb7e85107cbdd2fa386de Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 18:02:48 -0300 Subject: [PATCH 07/14] ADGStreamerApp: add GstBus monitoring Added `GstBus` monitoring for pipeline events such as state changes, errors, and End Of Stream (EOS), keeping the driver updated with the current pipeline execution state. --- ADGStreamerApp/src/ADGStreamer.cpp | 74 +++++++++++++++++++++++++----- ADGStreamerApp/src/ADGStreamer.h | 5 ++ 2 files changed, 67 insertions(+), 12 deletions(-) diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp index 518cd96..437bec5 100755 --- a/ADGStreamerApp/src/ADGStreamer.cpp +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -9,11 +9,9 @@ ADGStreamer::ADGStreamer(const char *portName, int maxBuffers, size_t maxMemory, { initializeGStreamer(); - epicsThreadCreate("ADGSTAcquire", - epicsThreadPriorityMedium, - epicsThreadGetStackSize(epicsThreadStackMedium), - (EPICSTHREADFUNC)acquisitionTaskC, - this); + epicsThreadCreate("ADGSTAcquire", epicsThreadPriorityMedium, + epicsThreadGetStackSize(epicsThreadStackMedium), + (EPICSTHREADFUNC)acquisitionTaskC, this); } ADGStreamer::~ADGStreamer() { stopPipeline(); } @@ -45,17 +43,13 @@ void ADGStreamer::acquisitionTask() epicsTimeStamp startTime, endTime; double elapsedTime, delay; - while (true) - { + while (true) { epicsTimeGetCurrent(&startTime); lock(); getIntegerParam(ADAcquire, &acquire); - if (acquire) - { - } - else - { + if (acquire) { + processBusMessage(); } unlock(); @@ -124,6 +118,11 @@ bool ADGStreamer::stopPipeline() sink_ = nullptr; } + if (bus_) { + gst_object_unref(bus_); + bus_ = nullptr; + } + if (pipeline_) { gst_object_unref(pipeline_); pipeline_ = nullptr; @@ -157,6 +156,11 @@ bool ADGStreamer::createPipeline() return false; } + bus_ = gst_element_get_bus(pipeline_); + if (!bus_) { + return false; + } + sink_ = gst_bin_get_by_name(GST_BIN(pipeline_), "appsink"); if (!sink_) { @@ -325,6 +329,52 @@ bool ADGStreamer::updateSample(GstCaps *caps) 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) { + g_error_free(error); + } + if (debug) { + g_free(debug); + } + break; + } + + case GST_MESSAGE_EOS: { + 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); + } + + 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) { diff --git a/ADGStreamerApp/src/ADGStreamer.h b/ADGStreamerApp/src/ADGStreamer.h index 291a64a..78ea37e 100755 --- a/ADGStreamerApp/src/ADGStreamer.h +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -36,14 +36,19 @@ class ADGStreamer : public ADDriver { 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 From d0df4695748ff7c14e079ece557f05ff7acbfa74 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Thu, 6 Aug 2026 11:20:55 -0300 Subject: [PATCH 08/14] ADGStreamerApp: introduce PipelineBuilder architecture for pipeline construction Introduce the PipelineBuilder architecture to separate GStreamer pipeline construction from the ADGStreamer driver implementation. Common pipeline stages, such as queue, decoder, converter, and sink, are implemented in PipelineBuilder, allowing them to be shared across different communication protocols. Protocol-specific implementations, such as PipelineBuilderRTSP, are responsible only for configuring protocol-related pipeline elements, making it easier to support additional protocols such as HTTP, RTMP, and SRT without requiring changes to the driver implementation. --- ADGStreamerApp/src/ADGStreamer.cpp | 97 +++++++++++++-- ADGStreamerApp/src/ADGStreamer.h | 6 + ADGStreamerApp/src/Makefile | 2 + ADGStreamerApp/src/PipelineBuilder.cpp | 92 ++++++++++++++ ADGStreamerApp/src/PipelineBuilder.h | 22 ++++ ADGStreamerApp/src/PipelineBuilderRTSP.cpp | 138 +++++++++++++++++++++ ADGStreamerApp/src/PipelineBuilderRTSP.h | 37 ++++++ ADGStreamerApp/src/PipelineTypes.h | 15 +++ 8 files changed, 398 insertions(+), 11 deletions(-) create mode 100755 ADGStreamerApp/src/PipelineBuilder.cpp create mode 100755 ADGStreamerApp/src/PipelineBuilder.h create mode 100755 ADGStreamerApp/src/PipelineBuilderRTSP.cpp create mode 100755 ADGStreamerApp/src/PipelineBuilderRTSP.h create mode 100755 ADGStreamerApp/src/PipelineTypes.h diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp index 437bec5..6836169 100755 --- a/ADGStreamerApp/src/ADGStreamer.cpp +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -14,7 +14,11 @@ ADGStreamer::ADGStreamer(const char *portName, int maxBuffers, size_t maxMemory, (EPICSTHREADFUNC)acquisitionTaskC, this); } -ADGStreamer::~ADGStreamer() { stopPipeline(); } +ADGStreamer::~ADGStreamer() +{ + stopPipeline(); + delete pipelineBuilder_; +} asynStatus ADGStreamer::writeInt32(asynUser *pasynUser, epicsInt32 value) { @@ -31,6 +35,19 @@ asynStatus ADGStreamer::writeInt32(asynUser *pasynUser, epicsInt32 value) return ADDriver::writeInt32(pasynUser, value); } +bool ADGStreamer::pipelineBuilder(PipelineBuilder *builder) +{ + if (!builder) + return false; + + if (pipelineBuilder_) + delete pipelineBuilder_; + + pipelineBuilder_ = builder; + + return true; +} + void ADGStreamer::acquisitionTaskC(void *drvPvt) { ADGStreamer *pPvt = static_cast(drvPvt); @@ -133,16 +150,11 @@ bool ADGStreamer::stopPipeline() bool ADGStreamer::createPipeline() { - std::string pipeline - = std::string("rtspsrc " - "location=rtsp://admin:r00tr00t@10.20.21.40:554/cam/" - "realmonitor?channel=1&subtype=0 protocols=udp latency=0 " - "drop-on-latency=true") - + std::string(" ! rtph264depay ! h264parse config-interval=-1") - + std::string( - " ! queue max-size-buffers=1 leaky=downstream silent=true") - + std::string(" ! avdec_h264 ! videoconvert ! video/x-raw,format=RGB ! " - "appsink name=appsink sync=false"); + if (!pipelineBuilder_) { + return false; + } + + std::string pipeline = pipelineBuilder_->build(); GError *error = nullptr; pipeline_ = gst_parse_launch(pipeline.c_str(), &error); @@ -401,9 +413,72 @@ static void ADGStreamerDriveCallFunc(const iocshArgBuf *args) 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.h b/ADGStreamerApp/src/ADGStreamer.h index 78ea37e..885e5e7 100755 --- a/ADGStreamerApp/src/ADGStreamer.h +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -9,6 +9,10 @@ #include #include +#include "PipelineBuilder.h" +#include "PipelineBuilderRTSP.h" +#include "PipelineTypes.h" + class ADGStreamer : public ADDriver { public: @@ -18,6 +22,8 @@ class ADGStreamer : public ADDriver { virtual asynStatus writeInt32(asynUser *pasynUser, epicsInt32 value); + bool pipelineBuilder(PipelineBuilder *builder); + protected: private: static void acquisitionTaskC(void *drvPvt); diff --git a/ADGStreamerApp/src/Makefile b/ADGStreamerApp/src/Makefile index d3af65f..1d28ad2 100644 --- a/ADGStreamerApp/src/Makefile +++ b/ADGStreamerApp/src/Makefile @@ -18,6 +18,8 @@ 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 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 From b1ce8e96fbb89f3b5644688842f9da2f6962a9e5 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Thu, 6 Aug 2026 12:21:00 -0300 Subject: [PATCH 09/14] ADGStreamerApp: include record read pipeline string --- ADGStreamerApp/Db/ADGStreamer.template | 11 +++++++++++ ADGStreamerApp/Db/Makefile | 1 + ADGStreamerApp/src/ADGStreamer.cpp | 5 +++++ ADGStreamerApp/src/ADGStreamer.h | 4 ++++ 4 files changed, 21 insertions(+) create mode 100755 ADGStreamerApp/Db/ADGStreamer.template 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 index 6836169..0b5e468 100755 --- a/ADGStreamerApp/src/ADGStreamer.cpp +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -7,6 +7,8 @@ ADGStreamer::ADGStreamer(const char *portName, int maxBuffers, size_t maxMemory, : ADDriver(portName, 1, 1, maxBuffers, maxMemory, 0, 0, ASYN_CANBLOCK, 1, priority, stackSize) { + createParam(ADGSTPipelineBuilder, asynParamOctet, &GST_PipelineBuilder); + initializeGStreamer(); epicsThreadCreate("ADGSTAcquire", epicsThreadPriorityMedium, @@ -156,6 +158,9 @@ bool ADGStreamer::createPipeline() std::string pipeline = pipelineBuilder_->build(); + setStringParam(GST_PipelineBuilder, pipeline.c_str()); + callParamCallbacks(); + GError *error = nullptr; pipeline_ = gst_parse_launch(pipeline.c_str(), &error); diff --git a/ADGStreamerApp/src/ADGStreamer.h b/ADGStreamerApp/src/ADGStreamer.h index 885e5e7..740a76b 100755 --- a/ADGStreamerApp/src/ADGStreamer.h +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -13,6 +13,8 @@ #include "PipelineBuilderRTSP.h" #include "PipelineTypes.h" +#define ADGSTPipelineBuilder "GST_PIPELINE_STRING" + class ADGStreamer : public ADDriver { public: @@ -25,6 +27,8 @@ class ADGStreamer : public ADDriver { bool pipelineBuilder(PipelineBuilder *builder); protected: + int GST_PipelineBuilder; + private: static void acquisitionTaskC(void *drvPvt); void acquisitionTask(); From eba3322e95ba42792743e0ec9f3f6aa87d5dadbe Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Thu, 6 Aug 2026 15:05:26 -0300 Subject: [PATCH 10/14] ADGStreamerApp: centralize driver status updates Add the reportStatus() helper method to centralize driver status updates and status message reporting. --- ADGStreamerApp/src/ADGStreamer.cpp | 41 +++++++++++++++++++++++++++++- ADGStreamerApp/src/ADGStreamer.h | 1 + 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/ADGStreamerApp/src/ADGStreamer.cpp b/ADGStreamerApp/src/ADGStreamer.cpp index 0b5e468..4b87d60 100755 --- a/ADGStreamerApp/src/ADGStreamer.cpp +++ b/ADGStreamerApp/src/ADGStreamer.cpp @@ -10,6 +10,7 @@ ADGStreamer::ADGStreamer(const char *portName, int maxBuffers, size_t maxMemory, createParam(ADGSTPipelineBuilder, asynParamOctet, &GST_PipelineBuilder); initializeGStreamer(); + reportStatus("GStreamer initialized", ADStatusIdle); epicsThreadCreate("ADGSTAcquire", epicsThreadPriorityMedium, epicsThreadGetStackSize(epicsThreadStackMedium), @@ -50,6 +51,13 @@ bool ADGStreamer::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); @@ -91,6 +99,8 @@ bool ADGStreamer::startPipeline() return false; } + reportStatus("Starting pipeline", ADStatusInitializing); + if (!createPipeline()) { return false; } @@ -101,6 +111,7 @@ bool ADGStreamer::startPipeline() ret = gst_element_set_state(pipeline_, GST_STATE_PLAYING); if (ret == GST_STATE_CHANGE_FAILURE) { + reportStatus("Could not start GStreamer pipeline", ADStatusError); return false; } @@ -108,19 +119,24 @@ bool ADGStreamer::startPipeline() 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; } @@ -128,6 +144,7 @@ bool ADGStreamer::stopPipeline() gst_element_get_state(pipeline_, &state, nullptr, GST_CLOCK_TIME_NONE); if (state != GST_STATE_NULL) { + reportStatus("Failed to stop pipeline", ADStatusError); return false; } } @@ -147,15 +164,20 @@ bool ADGStreamer::stopPipeline() 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()); @@ -165,22 +187,27 @@ bool ADGStreamer::createPipeline() pipeline_ = gst_parse_launch(pipeline.c_str(), &error); if (!pipeline_) { - if (error) { } + 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; } @@ -191,6 +218,8 @@ bool ADGStreamer::createPipeline() g_signal_connect( sink_, "new-sample", G_CALLBACK(onNewSampleCallback), this); + reportStatus("Pipeline created", ADStatusIdle); + return true; } @@ -360,6 +389,7 @@ bool ADGStreamer::processBusMessage() gchar *debug = nullptr; gst_message_parse_error(message, &error, &debug); if (error) { + reportStatus(error->message, ADStatusError); g_error_free(error); } if (debug) { @@ -369,6 +399,7 @@ bool ADGStreamer::processBusMessage() } case GST_MESSAGE_EOS: { + reportStatus("End of stream", ADStatusIdle); break; } @@ -379,6 +410,14 @@ bool ADGStreamer::processBusMessage() 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; diff --git a/ADGStreamerApp/src/ADGStreamer.h b/ADGStreamerApp/src/ADGStreamer.h index 740a76b..94dd25e 100755 --- a/ADGStreamerApp/src/ADGStreamer.h +++ b/ADGStreamerApp/src/ADGStreamer.h @@ -25,6 +25,7 @@ class ADGStreamer : public ADDriver { virtual asynStatus writeInt32(asynUser *pasynUser, epicsInt32 value); bool pipelineBuilder(PipelineBuilder *builder); + void reportStatus(const char *message, ADStatus_t status); protected: int GST_PipelineBuilder; From 38ad0f5266f827e2d4137dfebfd2b1775edb3e95 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 13:11:21 -0300 Subject: [PATCH 11/14] iocs: link ioc to support module --- iocs/ADGStreamerIOC/ADGStreamerApp/src/Makefile | 5 +++++ iocs/ADGStreamerIOC/configure/CONFIG_SITE | 8 ++++++++ iocs/ADGStreamerIOC/configure/RELEASE | 3 +++ 3 files changed, 16 insertions(+) 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 From f9f80105428e02416d97ed44f4a773183b2b84e4 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Wed, 5 Aug 2026 13:19:34 -0300 Subject: [PATCH 12/14] iocs: configure st.cmd --- .../iocBoot/iocADGStreamer/st.cmd | 73 ++++++++++++++++--- 1 file changed, 64 insertions(+), 9 deletions(-) mode change 100644 => 100755 iocs/ADGStreamerIOC/iocBoot/iocADGStreamer/st.cmd 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") From 3717679a7a03b964a480c22afdd47485b754b6c8 Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Thu, 6 Aug 2026 16:26:52 -0300 Subject: [PATCH 13/14] readme: include readme file --- README.md | 14 ++++++++++++++ 1 file changed, 14 insertions(+) create mode 100755 README.md 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 From 8368c6135c897bff7517451165b72eb5abfc02cb Mon Sep 17 00:00:00 2001 From: Guilherme Rodrigues de Lima Date: Thu, 6 Aug 2026 16:52:43 -0300 Subject: [PATCH 14/14] docs: add initial ADGStreamer documentation --- docs/ADGStreamer/ADGStreamer.rst | 151 +++++++++++++++++++++++++++++++ 1 file changed, 151 insertions(+) create mode 100644 docs/ADGStreamer/ADGStreamer.rst 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)