2024-03-12 21:13:13 +01:00
|
|
|
// #define DEBUGTRACE_ENABLED
|
|
|
|
#include <pybind11/pybind11.h>
|
|
|
|
#include <pybind11/pytypes.h>
|
|
|
|
|
|
|
|
#include <armadillo>
|
|
|
|
#include <atomic>
|
|
|
|
|
2022-10-11 09:43:36 +02:00
|
|
|
#include "arma_npy.h"
|
2022-10-06 21:13:21 +02:00
|
|
|
#include "debugtrace.hpp"
|
2023-01-27 14:56:46 +01:00
|
|
|
#include "lasp_clip.h"
|
2023-06-07 21:49:07 +02:00
|
|
|
#include "lasp_daq.h"
|
|
|
|
#include "lasp_daqdata.h"
|
|
|
|
#include "lasp_ppm.h"
|
2022-10-06 21:13:21 +02:00
|
|
|
#include "lasp_rtaps.h"
|
2023-01-05 10:35:47 +01:00
|
|
|
#include "lasp_rtsignalviewer.h"
|
2023-06-06 16:05:24 +02:00
|
|
|
#include "lasp_streammgr.h"
|
2023-06-07 21:49:07 +02:00
|
|
|
#include "lasp_threadedindatahandler.h"
|
2022-07-29 09:32:26 +02:00
|
|
|
|
|
|
|
using namespace std::literals::chrono_literals;
|
|
|
|
using std::cerr;
|
|
|
|
using std::endl;
|
2024-03-13 12:19:24 +01:00
|
|
|
using rte = std::runtime_error;
|
|
|
|
using Lck = std::scoped_lock<std::recursive_mutex>;
|
2022-07-29 09:32:26 +02:00
|
|
|
|
|
|
|
namespace py = pybind11;
|
|
|
|
|
|
|
|
/**
|
|
|
|
* @brief Generate a Numpy array from daqdata, does *NOT* create a copy of the
|
2022-08-01 17:26:22 +02:00
|
|
|
* data!. Instead, it shares the data from the DaqData container.
|
2022-07-29 09:32:26 +02:00
|
|
|
*
|
|
|
|
* @tparam T The type of the stored sample
|
|
|
|
* @param d The daqdata to convert
|
|
|
|
*
|
|
|
|
* @return Numpy array
|
|
|
|
*/
|
2022-10-11 09:43:36 +02:00
|
|
|
template <typename T, bool copy = false>
|
|
|
|
py::array_t<T> getPyArrayNoCpy(const DaqData &d) {
|
2022-07-29 09:32:26 +02:00
|
|
|
// https://github.com/pybind/pybind11/issues/323
|
|
|
|
//
|
|
|
|
// When a valid object is passed as 'base', it tells pybind not to take
|
|
|
|
// ownership of the data, because 'base' will own it. In fact 'packet' will
|
|
|
|
// own it, but - psss! - , we don't tell it to pybind... Alos note that ANY
|
2022-10-10 19:17:38 +02:00
|
|
|
// valid object is good for this purpose, so I choose "str"...
|
|
|
|
|
|
|
|
py::str dummyDataOwner;
|
|
|
|
/*
|
|
|
|
* Signature:
|
|
|
|
array_t(ShapeContainer shape,
|
|
|
|
StridesContainer strides,
|
|
|
|
const T *ptr = nullptr,
|
|
|
|
handle base = handle());
|
|
|
|
*/
|
|
|
|
|
|
|
|
return py::array_t<T>(
|
2024-03-12 21:13:13 +01:00
|
|
|
py::array::ShapeContainer({d.nframes, d.nchannels}), // Shape
|
2022-10-10 19:17:38 +02:00
|
|
|
|
2024-03-12 21:13:13 +01:00
|
|
|
py::array::StridesContainer( // Strides
|
2022-10-10 19:17:38 +02:00
|
|
|
{sizeof(T),
|
2024-03-12 21:13:13 +01:00
|
|
|
sizeof(T) * d.nframes}), // Strides (in bytes) for each index
|
2022-10-10 19:17:38 +02:00
|
|
|
|
|
|
|
reinterpret_cast<T *>(
|
2024-03-12 21:13:13 +01:00
|
|
|
const_cast<DaqData &>(d).raw_ptr()), // Pointer to buffer
|
2022-10-10 19:17:38 +02:00
|
|
|
|
2024-03-12 21:13:13 +01:00
|
|
|
dummyDataOwner // As stated above, now Numpy does not take ownership of
|
|
|
|
// the data pointer.
|
2022-10-10 19:17:38 +02:00
|
|
|
);
|
|
|
|
}
|
|
|
|
|
2022-10-11 09:43:36 +02:00
|
|
|
template <typename T, bool copy = false>
|
|
|
|
py::array_t<d> dmat_to_ndarray(const DaqData &d) {
|
2022-10-10 19:17:38 +02:00
|
|
|
// https://github.com/pybind/pybind11/issues/323
|
|
|
|
//
|
|
|
|
// When a valid object is passed as 'base', it tells pybind not to take
|
|
|
|
// ownership of the data, because 'base' will own it. In fact 'packet' will
|
|
|
|
// own it, but - psss! - , we don't tell it to pybind... Alos note that ANY
|
|
|
|
// valid object is good for this purpose, so I choose "str"...
|
2022-07-29 09:32:26 +02:00
|
|
|
|
|
|
|
py::str dummyDataOwner;
|
|
|
|
/*
|
|
|
|
* Signature:
|
|
|
|
array_t(ShapeContainer shape,
|
|
|
|
StridesContainer strides,
|
|
|
|
const T *ptr = nullptr,
|
|
|
|
handle base = handle());
|
|
|
|
*/
|
|
|
|
|
|
|
|
return py::array_t<T>(
|
2024-03-12 21:13:13 +01:00
|
|
|
py::array::ShapeContainer({d.nframes, d.nchannels}), // Shape
|
2022-08-14 21:00:22 +02:00
|
|
|
|
2024-03-12 21:13:13 +01:00
|
|
|
py::array::StridesContainer( // Strides
|
2022-08-14 21:00:22 +02:00
|
|
|
{sizeof(T),
|
2024-03-12 21:13:13 +01:00
|
|
|
sizeof(T) * d.nframes}), // Strides (in bytes) for each index
|
2022-08-01 17:26:22 +02:00
|
|
|
|
2022-08-14 21:00:22 +02:00
|
|
|
reinterpret_cast<T *>(
|
2024-03-12 21:13:13 +01:00
|
|
|
const_cast<DaqData &>(d).raw_ptr()), // Pointer to buffer
|
2022-08-01 17:26:22 +02:00
|
|
|
|
2024-03-12 21:13:13 +01:00
|
|
|
dummyDataOwner // As stated above, now Numpy does not take ownership of
|
|
|
|
// the data pointer.
|
2022-08-14 21:00:22 +02:00
|
|
|
);
|
2022-07-29 09:32:26 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
/**
|
2023-06-20 17:08:55 +02:00
|
|
|
* @brief Wraps the ThreadedInDataHandler such that it calls a Python callback
|
|
|
|
* with a buffer of sample data. Converts DaqData objects to Numpy arrays and
|
|
|
|
* calls Python given as argument to the constructor
|
2022-07-29 09:32:26 +02:00
|
|
|
*/
|
2023-06-09 10:43:04 +02:00
|
|
|
class PyIndataHandler : public ThreadedInDataHandler<PyIndataHandler> {
|
2022-07-29 09:32:26 +02:00
|
|
|
/**
|
2022-10-04 09:27:27 +02:00
|
|
|
* @brief The callback functions that is called.
|
2022-07-29 09:32:26 +02:00
|
|
|
*/
|
2024-03-13 12:19:24 +01:00
|
|
|
py::object _cb, _reset_callback;
|
2024-03-12 21:13:13 +01:00
|
|
|
std::atomic<bool> _done{false};
|
2024-03-13 12:19:24 +01:00
|
|
|
std::recursive_mutex _mtx;
|
2022-07-29 09:32:26 +02:00
|
|
|
|
2024-03-06 21:41:04 +01:00
|
|
|
public:
|
2023-06-09 10:43:04 +02:00
|
|
|
/**
|
|
|
|
* @brief Initialize PyIndataHandler
|
|
|
|
*
|
|
|
|
* @param mgr StreamMgr handle
|
|
|
|
* @param cb Python callback that is called with Numpy input data from device
|
|
|
|
* @param reset_callback Python callback that is called with a Daq pointer.
|
|
|
|
* Careful: do not store this handle, as it is only valid as long as reset()
|
|
|
|
* is called, when a stream stops, this pointer / handle will dangle.
|
|
|
|
*/
|
2023-06-07 21:49:07 +02:00
|
|
|
PyIndataHandler(SmgrHandle mgr, py::function cb, py::function reset_callback)
|
2024-03-13 12:19:24 +01:00
|
|
|
: ThreadedInDataHandler(mgr),
|
|
|
|
_cb(py::weakref(cb)),
|
|
|
|
_reset_callback(py::weakref(reset_callback)) {
|
2022-08-14 21:00:22 +02:00
|
|
|
DEBUGTRACE_ENTER;
|
2024-03-13 12:19:24 +01:00
|
|
|
// cerr << "Thread ID: " << std::this_thread::get_id() << endl;
|
2022-10-10 19:17:38 +02:00
|
|
|
/// Start should be called externally, as at constructor time no virtual
|
|
|
|
/// functions should be called.
|
2024-03-13 12:19:24 +01:00
|
|
|
if (_cb().is_none() || _reset_callback().is_none()) {
|
|
|
|
throw rte("cb or reset_callback is none!");
|
|
|
|
}
|
2023-06-06 16:05:24 +02:00
|
|
|
startThread();
|
2022-08-14 21:00:22 +02:00
|
|
|
}
|
2022-07-29 09:32:26 +02:00
|
|
|
~PyIndataHandler() {
|
|
|
|
DEBUGTRACE_ENTER;
|
2024-03-13 12:19:24 +01:00
|
|
|
// cerr << "Thread ID: " << std::this_thread::get_id() << endl;
|
2023-06-20 17:08:55 +02:00
|
|
|
/// Callback cannot be called, which results in a deadlock on the GIL
|
|
|
|
/// without this release.
|
|
|
|
py::gil_scoped_release release;
|
2024-03-13 12:19:24 +01:00
|
|
|
DEBUGTRACE_PRINT("Gil released");
|
|
|
|
_done = true;
|
2023-06-06 16:05:24 +02:00
|
|
|
stopThread();
|
2022-07-29 09:32:26 +02:00
|
|
|
}
|
2023-06-09 10:43:04 +02:00
|
|
|
/**
|
|
|
|
* @brief Calls the reset callback in Python.
|
|
|
|
*
|
|
|
|
* @param daq Daq device, or nullptr in case no input stream is running.
|
|
|
|
*/
|
2024-03-12 21:13:13 +01:00
|
|
|
void reset(const Daq *daqi) {
|
2023-06-07 21:49:07 +02:00
|
|
|
DEBUGTRACE_ENTER;
|
2024-03-13 12:19:24 +01:00
|
|
|
// cerr << "Thread ID: " << std::this_thread::get_id() << endl;
|
2024-03-06 21:41:04 +01:00
|
|
|
if (_done) return;
|
2024-03-12 21:13:13 +01:00
|
|
|
{
|
|
|
|
try {
|
|
|
|
py::object reset_callback = _reset_callback();
|
|
|
|
if (reset_callback.is_none()) {
|
2024-03-13 12:19:24 +01:00
|
|
|
DEBUGTRACE_PRINT("reset_callback is none, weakref killed");
|
2024-03-12 21:13:13 +01:00
|
|
|
_done = true;
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
if (daqi != nullptr) {
|
|
|
|
assert(reset_callback);
|
|
|
|
reset_callback(daqi);
|
|
|
|
} else {
|
|
|
|
assert(reset_callback);
|
|
|
|
reset_callback(py::none());
|
|
|
|
}
|
|
|
|
} catch (py::error_already_set &e) {
|
|
|
|
cerr << "*************** Error calling reset callback!\n";
|
|
|
|
cerr << e.what() << endl;
|
|
|
|
cerr << "*************** \n";
|
|
|
|
/// Throwing a runtime error here does not work out one way or another.
|
|
|
|
/// Therefore, it is better to dive out and prevent undefined behaviour
|
|
|
|
abort();
|
|
|
|
/* throw std::runtime_error(e.what()); */
|
|
|
|
} catch (std::exception &e) {
|
|
|
|
cerr << "Caught unknown exception in reset callback:" << e.what()
|
|
|
|
<< endl;
|
|
|
|
abort();
|
2023-06-07 21:49:07 +02:00
|
|
|
}
|
2024-03-13 12:19:24 +01:00
|
|
|
} // end of GIL scope
|
|
|
|
} // end of function reset()
|
2022-07-29 09:32:26 +02:00
|
|
|
|
|
|
|
/**
|
2023-06-09 10:43:04 +02:00
|
|
|
* @brief Calls the Python callback method / function with a Numpy array of
|
|
|
|
* stream data.
|
2022-07-29 09:32:26 +02:00
|
|
|
*/
|
2023-06-09 10:43:04 +02:00
|
|
|
void inCallback(const DaqData &d) {
|
2024-03-13 12:19:24 +01:00
|
|
|
DEBUGTRACE_ENTER;
|
|
|
|
// cerr << "=== Enter incallback for thread ID: " << std::this_thread::get_id() << endl;
|
2022-07-29 09:32:26 +02:00
|
|
|
|
|
|
|
using DataType = DataTypeDescriptor::DataType;
|
2024-03-12 21:13:13 +01:00
|
|
|
if (_done) {
|
|
|
|
DEBUGTRACE_PRINT("Early stop, done");
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
{
|
2024-03-13 12:19:24 +01:00
|
|
|
DEBUGTRACE_PRINT("================ TRYING TO OBTAIN GIL in inCallback...");
|
2023-06-20 17:08:55 +02:00
|
|
|
py::gil_scoped_acquire acquire;
|
2024-03-12 21:13:13 +01:00
|
|
|
try {
|
|
|
|
py::object py_bool;
|
|
|
|
py::object cb = _cb();
|
|
|
|
if (cb.is_none()) {
|
|
|
|
DEBUGTRACE_PRINT("cb is none, weakref killed");
|
|
|
|
_done = true;
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
switch (d.dtype) {
|
|
|
|
case (DataType::dtype_int8): {
|
|
|
|
py_bool = cb(getPyArrayNoCpy<int8_t>(d));
|
|
|
|
} break;
|
|
|
|
case (DataType::dtype_int16): {
|
|
|
|
py_bool = cb(getPyArrayNoCpy<int16_t>(d));
|
|
|
|
} break;
|
|
|
|
case (DataType::dtype_int32): {
|
|
|
|
py_bool = cb(getPyArrayNoCpy<int32_t>(d));
|
|
|
|
} break;
|
|
|
|
case (DataType::dtype_fl32): {
|
|
|
|
py_bool = cb(getPyArrayNoCpy<float>(d));
|
|
|
|
} break;
|
|
|
|
case (DataType::dtype_fl64): {
|
|
|
|
py_bool = cb(getPyArrayNoCpy<double>(d));
|
|
|
|
} break;
|
|
|
|
default:
|
|
|
|
throw std::runtime_error("BUG");
|
|
|
|
} // End of switch
|
|
|
|
|
|
|
|
bool res = py_bool.cast<bool>();
|
|
|
|
if (res == false) {
|
|
|
|
DEBUGTRACE_PRINT("Setting callbacks to None")
|
|
|
|
_done = true;
|
|
|
|
|
|
|
|
// By doing this, we remove the references, but in the mean time this
|
|
|
|
// might also trigger removing Python objects. Including itself, as
|
|
|
|
// there is no reference to it anymore. The consequence is that the
|
|
|
|
// current object might be destroyed from this thread. However, if we
|
|
|
|
// do not remove these references and in lasp_record.py finish() is
|
|
|
|
// not called, we end up with not-garbage collected recordings in
|
|
|
|
// memory. This is also not good. How can we force Python to not yet
|
|
|
|
// destroy this object?
|
|
|
|
// cb.reset();
|
|
|
|
// reset_callback.reset();
|
|
|
|
}
|
|
|
|
} catch (py::error_already_set &e) {
|
2024-03-13 12:19:24 +01:00
|
|
|
cerr << "ERROR (BUG): Python raised exception from callback function: ";
|
2024-03-12 21:13:13 +01:00
|
|
|
cerr << e.what() << endl;
|
|
|
|
abort();
|
|
|
|
} catch (py::cast_error &e) {
|
|
|
|
cerr << e.what() << endl;
|
2024-03-13 12:19:24 +01:00
|
|
|
cerr << "ERROR (BUG): Python callback does not return boolean value."
|
|
|
|
<< endl;
|
2024-03-12 21:13:13 +01:00
|
|
|
abort();
|
|
|
|
} catch (std::exception &e) {
|
|
|
|
cerr << "Caught unknown exception in Python callback:" << e.what()
|
|
|
|
<< endl;
|
|
|
|
abort();
|
2024-03-06 21:41:04 +01:00
|
|
|
}
|
2024-03-12 21:13:13 +01:00
|
|
|
|
|
|
|
} // End of scope in which the GIL is acquired
|
2024-03-13 12:19:24 +01:00
|
|
|
// cerr << "=== LEAVE incallback for thread ID: " << std::this_thread::get_id() << endl;
|
|
|
|
|
|
|
|
} // End of function inCallback()
|
2022-07-29 09:32:26 +02:00
|
|
|
};
|
2022-10-06 21:13:21 +02:00
|
|
|
|
2022-07-29 09:32:26 +02:00
|
|
|
void init_datahandler(py::module &m) {
|
2022-10-10 19:17:38 +02:00
|
|
|
/// The C++ class is PyIndataHandler, but for Python, it is called
|
|
|
|
/// InDataHandler
|
2023-06-06 16:05:24 +02:00
|
|
|
py::class_<PyIndataHandler> pyidh(m, "InDataHandler");
|
2023-06-07 21:49:07 +02:00
|
|
|
pyidh.def(py::init<SmgrHandle, py::function, py::function>());
|
2022-07-29 09:32:26 +02:00
|
|
|
|
2022-10-06 21:13:21 +02:00
|
|
|
/// Peak Programme Meter
|
2023-06-06 16:05:24 +02:00
|
|
|
py::class_<PPMHandler> ppm(m, "PPMHandler");
|
2023-06-07 21:49:07 +02:00
|
|
|
ppm.def(py::init<SmgrHandle, const d>());
|
|
|
|
ppm.def(py::init<SmgrHandle>());
|
2022-10-04 09:27:27 +02:00
|
|
|
|
|
|
|
ppm.def("getCurrentValue", [](const PPMHandler &ppm) {
|
2023-06-20 17:16:56 +02:00
|
|
|
std::tuple<vd, arma::uvec> tp;
|
|
|
|
{
|
|
|
|
py::gil_scoped_release release;
|
|
|
|
tp = ppm.getCurrentValue();
|
|
|
|
}
|
2022-10-04 09:27:27 +02:00
|
|
|
|
2022-10-11 09:43:36 +02:00
|
|
|
return py::make_tuple(ColToNpy<d>(std::get<0>(tp)),
|
2023-01-05 10:35:47 +01:00
|
|
|
ColToNpy<arma::uword>(std::get<1>(tp)));
|
2022-10-06 21:13:21 +02:00
|
|
|
});
|
|
|
|
|
2023-01-27 14:56:46 +01:00
|
|
|
/// Clip Detector
|
2023-06-06 16:05:24 +02:00
|
|
|
py::class_<ClipHandler> clip(m, "ClipHandler");
|
2023-06-07 21:49:07 +02:00
|
|
|
clip.def(py::init<SmgrHandle>());
|
2023-01-27 14:56:46 +01:00
|
|
|
|
|
|
|
clip.def("getCurrentValue", [](const ClipHandler &clip) {
|
2023-06-20 17:16:56 +02:00
|
|
|
arma::uvec cval;
|
|
|
|
{
|
|
|
|
py::gil_scoped_release release;
|
|
|
|
cval = clip.getCurrentValue();
|
|
|
|
}
|
2023-01-27 14:56:46 +01:00
|
|
|
|
2024-03-12 21:13:13 +01:00
|
|
|
return ColToNpy<arma::uword>(cval); // something goes wrong here
|
2023-01-27 14:56:46 +01:00
|
|
|
});
|
|
|
|
|
2022-10-06 21:13:21 +02:00
|
|
|
/// Real time Aps
|
|
|
|
///
|
2023-06-06 16:05:24 +02:00
|
|
|
py::class_<RtAps> rtaps(m, "RtAps");
|
2024-03-12 21:13:13 +01:00
|
|
|
rtaps.def(py::init<SmgrHandle, // StreamMgr
|
|
|
|
Filter *const, // FreqWeighting filter
|
|
|
|
const us, // Nfft
|
|
|
|
const Window::WindowType, // Window
|
|
|
|
const d, // Overlap percentage 0<=o<100
|
2022-10-06 21:13:21 +02:00
|
|
|
|
2024-03-12 21:13:13 +01:00
|
|
|
const d // Time constant
|
2022-10-06 21:13:21 +02:00
|
|
|
>(),
|
2024-03-12 21:13:13 +01:00
|
|
|
py::arg("streammgr"), // StreamMgr
|
2022-10-06 21:13:21 +02:00
|
|
|
py::arg("preFilter").none(true),
|
|
|
|
/// Below list of arguments *SHOULD* be same as for
|
|
|
|
|
|
|
|
/// AvPowerSpectra constructor!
|
2024-03-12 21:13:13 +01:00
|
|
|
py::arg("nfft") = 2048, //
|
|
|
|
py::arg("windowType") = Window::WindowType::Hann, //
|
|
|
|
py::arg("overlap_percentage") = 50.0, //
|
|
|
|
py::arg("time_constant") = -1 //
|
2022-10-06 21:13:21 +02:00
|
|
|
);
|
|
|
|
|
|
|
|
rtaps.def("getCurrentValue", [](RtAps &rt) {
|
2023-06-20 17:16:56 +02:00
|
|
|
ccube val;
|
|
|
|
{
|
|
|
|
py::gil_scoped_release release;
|
|
|
|
val = rt.getCurrentValue();
|
|
|
|
}
|
2023-01-05 10:35:47 +01:00
|
|
|
return CubeToNpy<c>(val);
|
|
|
|
});
|
|
|
|
|
|
|
|
/// Real time Signal Viewer
|
|
|
|
///
|
2023-06-06 16:05:24 +02:00
|
|
|
py::class_<RtSignalViewer> rtsv(m, "RtSignalViewer");
|
2024-03-12 21:13:13 +01:00
|
|
|
rtsv.def(py::init<SmgrHandle, // StreamMgr
|
|
|
|
const d, // Time history
|
|
|
|
const us, // Resolution
|
|
|
|
const us // Channel number
|
2023-01-05 10:35:47 +01:00
|
|
|
>());
|
|
|
|
|
|
|
|
rtsv.def("getCurrentValue", [](RtSignalViewer &rt) {
|
2023-06-20 17:16:56 +02:00
|
|
|
dmat val;
|
|
|
|
{
|
|
|
|
py::gil_scoped_release release;
|
|
|
|
val = rt.getCurrentValue();
|
|
|
|
}
|
2023-01-05 10:35:47 +01:00
|
|
|
return MatToNpy<d>(val);
|
2022-10-04 09:27:27 +02:00
|
|
|
});
|
2022-07-29 09:32:26 +02:00
|
|
|
}
|