Skip to content

Commit 268c4f5

Browse files
committed
Update DelegateMQ library
1 parent d0b6965 commit 268c4f5

42 files changed

Lines changed: 1624 additions & 286 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎DelegateMQ/DelegateMQ.h‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,7 @@
8484
// Valid for any platform where a Mutex is defined in DelegateOpt.h
8585
#if defined(DMQ_THREAD_STDLIB) || \
8686
defined(DMQ_THREAD_WIN32) || \
87+
defined(DMQ_THREAD_POSIX) || \
8788
defined(DMQ_THREAD_FREERTOS) || \
8889
defined(DMQ_THREAD_THREADX) || \
8990
defined(DMQ_THREAD_ZEPHYR) || \
@@ -104,6 +105,7 @@
104105
// Valid for any platform that implements the IThread interface
105106
#if defined(DMQ_THREAD_STDLIB) || \
106107
defined(DMQ_THREAD_WIN32) || \
108+
defined(DMQ_THREAD_POSIX) || \
107109
defined(DMQ_THREAD_FREERTOS) || \
108110
defined(DMQ_THREAD_THREADX) || \
109111
defined(DMQ_THREAD_ZEPHYR) || \
@@ -133,6 +135,9 @@
133135
#elif defined(DMQ_THREAD_WIN32)
134136
#include "port/os/win32/Win32Thread.h"
135137
#include "port/os/common/ThreadMsg.h"
138+
#elif defined(DMQ_THREAD_POSIX)
139+
#include "port/os/posix/PosixThread.h"
140+
#include "port/os/common/ThreadMsg.h"
136141
#elif defined(DMQ_THREAD_FREERTOS)
137142
#include "port/os/freertos/FreeRTOSThread.h"
138143
#include "port/os/common/ThreadMsg.h"

‎DelegateMQ/Port.cmake‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,8 +9,14 @@ if (DMQ_THREAD STREQUAL "DMQ_THREAD_STDLIB")
99
elseif (DMQ_THREAD STREQUAL "DMQ_THREAD_WIN32")
1010
add_compile_definitions(DMQ_THREAD_WIN32)
1111
file(GLOB THREAD_SOURCES CONFIGURE_DEPENDS
12-
"${DMQ_ROOT_DIR}/port/os/win32/*.c*"
13-
"${DMQ_ROOT_DIR}/port/os/win32/*.h"
12+
"${DMQ_ROOT_DIR}/port/os/win32/*.c*"
13+
"${DMQ_ROOT_DIR}/port/os/win32/*.h"
14+
)
15+
elseif (DMQ_THREAD STREQUAL "DMQ_THREAD_POSIX")
16+
add_compile_definitions(DMQ_THREAD_POSIX)
17+
file(GLOB THREAD_SOURCES CONFIGURE_DEPENDS
18+
"${DMQ_ROOT_DIR}/port/os/posix/*.c*"
19+
"${DMQ_ROOT_DIR}/port/os/posix/*.h"
1420
)
1521
elseif (DMQ_THREAD STREQUAL "DMQ_THREAD_FREERTOS")
1622
add_compile_definitions(DMQ_THREAD_FREERTOS)

‎DelegateMQ/delegate/DelegateMQConfig_Default.h‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,17 @@
2525
#define DMQ_DEFAULT_QUEUE_SIZE 20
2626
#endif
2727

28+
// Fallback queue size (maxQueueSize == 0) for the desktop stdlib/Win32 Thread ports
29+
// only. Unlike DMQ_DEFAULT_QUEUE_SIZE (sized for real RTOS RAM constraints), these
30+
// ports back their queue with a plain std::deque, so maxQueueSize == 0 otherwise
31+
// means genuinely unbounded growth if the destination thread is dead/stuck. This
32+
// is a high-water-mark safety net, not a throughput limiter -- large enough to
33+
// never interfere with normal desktop bursts, but bounded so a truly dead consumer
34+
// hits a clear FullPolicy::FAULT instead of growing the heap without limit.
35+
#ifndef DMQ_THREAD_DESKTOP_QUEUE_SIZE
36+
#define DMQ_THREAD_DESKTOP_QUEUE_SIZE 1000
37+
#endif
38+
2839
#ifndef DMQ_MAX_WATCHDOG_THREADS
2940
#define DMQ_MAX_WATCHDOG_THREADS 16
3041
#endif

‎DelegateMQ/delegate/DelegateMQConfig_Template.h‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,9 +20,17 @@
2020
/// Signals with <= this many subscribers are invoked without heap allocation.
2121
#define DMQ_SIGNAL_SBO_COUNT 8
2222

23-
/// Default internal message queue depth for all dmq::os::Thread ports.
23+
/// Default internal message queue depth (maxQueueSize == 0) for the RTOS
24+
/// dmq::os::Thread ports (FreeRTOS, Zephyr, ThreadX, CMSIS-RTOS2, NuttX), where
25+
/// the backing queue primitive requires a fixed capacity at creation.
2426
#define DMQ_DEFAULT_QUEUE_SIZE 20
2527

28+
/// Fallback queue size (maxQueueSize == 0) for the desktop stdlib/Win32 Thread
29+
/// ports only. A high-water-mark safety net (not a throughput limiter) against
30+
/// unbounded growth if the destination thread is dead/stuck -- large enough to
31+
/// never interfere with normal desktop bursts.
32+
#define DMQ_THREAD_DESKTOP_QUEUE_SIZE 1000
33+
2634
/// Max number of threads that can be registered with the watchdog.
2735
#define DMQ_MAX_WATCHDOG_THREADS 16
2836

‎DelegateMQ/delegate/DelegateOpt.h‎

Lines changed: 33 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@
2929

3030
// --- PLATFORM AUTO-DETECTION ---
3131
// If no threading model is defined, attempt to auto-select a default
32-
#if !defined(DMQ_THREAD_STDLIB) && !defined(DMQ_THREAD_WIN32) && \
32+
#if !defined(DMQ_THREAD_STDLIB) && !defined(DMQ_THREAD_WIN32) && !defined(DMQ_THREAD_POSIX) && \
3333
!defined(DMQ_THREAD_FREERTOS) && !defined(DMQ_THREAD_THREADX) && \
3434
!defined(DMQ_THREAD_ZEPHYR) && !defined(DMQ_THREAD_CMSIS_RTOS2) && \
3535
!defined(DMQ_THREAD_NUTTX) && \
@@ -72,9 +72,10 @@
7272
// True when a real thread model is configured (desktop or embedded RTOS),
7373
// as opposed to DMQ_THREAD_NONE or no thread model at all (bare metal,
7474
// single-threaded). Named once here instead of hand-copying this same
75-
// 7-macro list at every call site that needs to know whether Mutex/
75+
// 8-macro list at every call site that needs to know whether Mutex/
7676
// ConditionVariable/std::thread-equivalent support exists.
7777
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT) || \
78+
defined(DMQ_THREAD_POSIX) || \
7879
defined(DMQ_THREAD_FREERTOS) || defined(DMQ_THREAD_THREADX) || \
7980
defined(DMQ_THREAD_ZEPHYR) || defined(DMQ_THREAD_CMSIS_RTOS2) || \
8081
defined(DMQ_THREAD_NUTTX)
@@ -168,8 +169,10 @@
168169
// later #include "extras/util/Fault.h" below is a harmless no-op.
169170
#include "extras/util/Fault.h"
170171

171-
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT)
172-
// Windows / Linux / macOS / Qt (Standard Library)
172+
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT) || defined(DMQ_THREAD_POSIX)
173+
// Windows / Linux / macOS / Qt (Standard Library) / POSIX (raw pthreads,
174+
// but Mutex/ConditionVariable/Clock/ThisThread still reuse std:: here --
175+
// see port/os/posix/PosixThread.h, the only file this port adds)
173176
#include <condition_variable>
174177
#include <thread>
175178
#elif defined(DMQ_THREAD_FREERTOS)
@@ -249,8 +252,8 @@ namespace dmq
249252
// @TODO: Change aliases to switch clock type globally if necessary
250253

251254
// --- CLOCK SELECTION ---
252-
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT)
253-
// Windows / Linux / macOS / Qt
255+
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT) || defined(DMQ_THREAD_POSIX)
256+
// Windows / Linux / macOS / Qt / POSIX
254257
using Clock = std::chrono::steady_clock;
255258

256259
#elif defined(DMQ_THREAD_FREERTOS)
@@ -296,8 +299,8 @@ namespace dmq
296299
// std;` together, and a same-named nested namespace would make unqualified
297300
// this_thread::sleep_for() calls in that code ambiguous against
298301
// std::this_thread.
299-
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT)
300-
// Windows / Linux / macOS / Qt -- std::this_thread is already portable here.
302+
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT) || defined(DMQ_THREAD_POSIX)
303+
// Windows / Linux / macOS / Qt / POSIX -- std::this_thread is already portable here.
301304
struct ThisThread {
302305
template<typename Rep, typename Period>
303306
static void sleep_for(std::chrono::duration<Rep, Period> d) { std::this_thread::sleep_for(d); }
@@ -362,7 +365,10 @@ namespace dmq
362365
#endif
363366

364367
/// @brief Policy applied when a dmq::os::Thread port's message queue is full.
365-
/// @details Only meaningful when the port's maxQueueSize > 0.
368+
/// @details Always enforced -- a constructor's maxQueueSize == 0 is a sentinel meaning
369+
/// "use this port's default capacity" (dmq::DEFAULT_QUEUE_SIZE for the RTOS ports,
370+
/// dmq::THREAD_DESKTOP_QUEUE_SIZE for stdlib/Win32), not "disable the cap." Every port's
371+
/// effective queue capacity is therefore always > 0 at runtime.
366372
/// - DROP: DispatchDelegate() silently discards the message and returns immediately.
367373
/// - FAULT: DispatchDelegate() triggers a system fault if the queue is full.
368374
/// - TIMEOUT: DispatchDelegate() waits up to dispatchTimeout, then logs and drops.
@@ -389,10 +395,20 @@ namespace dmq
389395
/// Override via DMQ_SIGNAL_SBO_COUNT in delegatemqconfig.h.
390396
inline constexpr size_t SIGNAL_SBO_COUNT = DMQ_SIGNAL_SBO_COUNT;
391397

392-
/// @brief Default internal queue size for all dmq::os::Thread ports.
398+
/// @brief Default internal queue size (maxQueueSize == 0) for the RTOS
399+
/// dmq::os::Thread ports, where the backing queue primitive requires a
400+
/// fixed capacity at creation.
393401
/// Override via DMQ_DEFAULT_QUEUE_SIZE in delegatemqconfig.h.
394402
inline constexpr size_t DEFAULT_QUEUE_SIZE = DMQ_DEFAULT_QUEUE_SIZE;
395403

404+
/// @brief Fallback queue size (maxQueueSize == 0) for the desktop
405+
/// stdlib/Win32 Thread ports only, which back their queue with a plain
406+
/// std::deque and would otherwise grow without bound if the destination
407+
/// thread is dead/stuck. A high-water-mark safety net, not a throughput
408+
/// limiter -- large enough to never interfere with normal desktop bursts.
409+
/// Override via DMQ_THREAD_DESKTOP_QUEUE_SIZE in delegatemqconfig.h.
410+
inline constexpr size_t THREAD_DESKTOP_QUEUE_SIZE = DMQ_THREAD_DESKTOP_QUEUE_SIZE;
411+
396412
/// @brief Max number of threads that can be monitored by the watchdog.
397413
/// Override via DMQ_MAX_WATCHDOG_THREADS in delegatemqconfig.h.
398414
inline constexpr size_t MAX_WATCHDOG_THREADS = DMQ_MAX_WATCHDOG_THREADS;
@@ -430,8 +446,8 @@ namespace dmq
430446
inline constexpr size_t MAX_TRANSPORT_MONITOR_PENDING = DMQ_TRANSPORT_MONITOR_MAX_PENDING;
431447

432448
// --- MUTEX / LOCK SELECTION ---
433-
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT)
434-
// Windows / Linux / macOS / Qt
449+
#if defined(DMQ_THREAD_STDLIB) || defined(DMQ_THREAD_WIN32) || defined(DMQ_THREAD_QT) || defined(DMQ_THREAD_POSIX)
450+
// Windows / Linux / macOS / Qt / POSIX
435451
using Mutex = std::mutex;
436452
using RecursiveMutex = std::recursive_mutex;
437453
// No ISR concept reachable from userspace on desktop OSes, and
@@ -482,15 +498,13 @@ namespace dmq
482498
using Mutex = dmq::os::ZephyrMutex;
483499
using RecursiveMutex = dmq::os::ZephyrRecursiveMutex;
484500
// ISR-safe (irq_lock()/irq_unlock(key), not a k_mutex) -- see
485-
// ZephyrCriticalSection.h. UNVERIFIED: no Zephyr SDK/west workspace is
486-
// available in this development environment to build and run it; review
487-
// before relying on it in production.
501+
// ZephyrCriticalSection.h.
488502
using CriticalSection = dmq::os::ZephyrCriticalSection;
489503
// No dmq::ConditionVariable port for Zephyr (no DMQ_HAS_CV), but
490504
// dmq::Semaphore is available via Zephyr's own native k_sem instead of
491505
// the generic condvar+mutex implementation -- see ZephyrSemaphore.h and
492506
// delegate/Semaphore.h. This is what makes DelegateAsyncWait available
493-
// here. UNVERIFIED, same caveat as CriticalSection above.
507+
// here.
494508
using Semaphore = dmq::os::ZephyrSemaphore;
495509
#define DMQ_HAS_SEMAPHORE
496510
template<typename T> using LockGuard = PortableLockGuard<T>;
@@ -508,8 +522,7 @@ namespace dmq
508522
// native condvar primitive to build one from), but dmq::Semaphore is
509523
// available via osSemaphore directly instead -- see
510524
// CmsisRtos2Semaphore.h and delegate/Semaphore.h. This is what makes
511-
// DelegateAsyncWait available here. UNVERIFIED, same caveat as
512-
// CriticalSection above.
525+
// DelegateAsyncWait available here.
513526
using Semaphore = dmq::os::CmsisRtos2Semaphore;
514527
#define DMQ_HAS_SEMAPHORE
515528
template<typename T> using LockGuard = PortableLockGuard<T>;
@@ -523,16 +536,14 @@ namespace dmq
523536
// ISR-safe (up_irq_save()/up_irq_restore(), NuttX's own architecture-
524537
// portable interrupt-masking primitive, not a pthread_mutex_t) -- see
525538
// NuttXCriticalSection.h, including its FLAT-vs-PROTECTED/KERNEL-build
526-
// caveat. UNVERIFIED: no NuttX toolchain/simulator is available in this
527-
// development environment to build and run it; review before relying
528-
// on it in production.
539+
// caveat.
529540
using CriticalSection = dmq::os::NuttXCriticalSection;
530541
// No dmq::ConditionVariable port for NuttX (no DMQ_HAS_CV) -- not
531542
// because NuttX lacks pthread_cond_t (it has a real one), but because
532543
// dmq::Semaphore is available via NuttX's own native sem_t instead of
533544
// the generic condvar+mutex implementation -- see NuttXSemaphore.h and
534545
// delegate/Semaphore.h. This is what makes DelegateAsyncWait available
535-
// here. UNVERIFIED, same caveat as CriticalSection above.
546+
// here.
536547
using Semaphore = dmq::os::NuttXSemaphore;
537548
#define DMQ_HAS_SEMAPHORE
538549
template<typename T> using LockGuard = PortableLockGuard<T>;

0 commit comments

Comments
 (0)