diff --git a/RBcore/msg_center.c b/RBcore/msg_center.c index 36b34c9..ad25385 100644 --- a/RBcore/msg_center.c +++ b/RBcore/msg_center.c @@ -45,6 +45,11 @@ typedef struct { MsgModule_t stModule; rd_ringbuf_t *pstRingBuf; // 消息队列(ringbuffer) +#if defined(USE_FREERTOS) + rd_mutex_t stMutex; // FreeRTOS互斥锁(任务上下文保护ringbuffer) +#endif + volatile uint8_t bSpinLock; // 自旋锁(ISR安全,始终可用) + uint32_t uiPrimask; // 关中断前的PRIMASK快照(裸机用) rd_sem_t stSem; // 信号量(阻塞等待) } MsgModuleEntry_t; @@ -68,6 +73,28 @@ static uint32_t g_uiNextModuleID = 1; // 从1开始,0表示无效 ((uiModuleID) > 0 && (uiModuleID) <= g_uiModuleCount ? \ (uiModuleID) - 1 : -1) +// 写入保护宏:FreeRTOS 用 mutex,裸机用关中断+自旋锁 +#if defined(USE_FREERTOS) +#define MSG_CENTER_LOCK(pEntry) Rd_MutexLock(&(pEntry)->stMutex, __func__, NULL) +#define MSG_CENTER_UNLOCK(pEntry) Rd_MutexUnlock(&(pEntry)->stMutex, __func__, NULL) +#define MSG_CENTER_LOCK_INIT(pEntry) ({ \ + (pEntry)->bSpinLock = 0; \ + Rd_MutexInit(&(pEntry)->stMutex); \ +}) +#define MSG_CENTER_LOCK_DESTROY(pEntry) Rd_MutexDestroy(&(pEntry)->stMutex) +#else +#define MSG_CENTER_LOCK(pEntry) do { \ + (pEntry)->uiPrimask = rd_irq_save(); \ + while (ATOMIC_EXCH(&(pEntry)->bSpinLock, 1, uint8_t, ATOMIC_ORDER_ACQUIRE)) {} \ + } while(0) +#define MSG_CENTER_UNLOCK(pEntry) do { \ + ATOMIC_STORE(&(pEntry)->bSpinLock, 0, uint8_t, ATOMIC_ORDER_RELEASE); \ + rd_irq_restore((pEntry)->uiPrimask); \ + } while(0) +#define MSG_CENTER_LOCK_INIT(pEntry) ({ (pEntry)->bSpinLock = 0; RD_SUCCESS; }) +#define MSG_CENTER_LOCK_DESTROY(pEntry) ((void)(pEntry)) +#endif + /***************************************************************************** 函 数 名 : MsgCenter_Init 功能描述 : 初始化消息中心 @@ -128,9 +155,17 @@ uint32_t MsgCenter_Register(const char *pszName, MsgHandler_t pfHandler) return 0; } + // 初始化写入锁 + if (MSG_CENTER_LOCK_INIT(&g_astModules[uiIndex]) != RD_SUCCESS) + { + rd_RingbufferDestroy(g_astModules[uiIndex].pstRingBuf); + return 0; + } + // 初始化信号量(初始值为0,表示无消息) if (Rd_SemInit(&g_astModules[uiIndex].stSem, 0, 0) != RD_SUCCESS) { + MSG_CENTER_LOCK_DESTROY(&g_astModules[uiIndex]); rd_RingbufferDestroy(g_astModules[uiIndex].pstRingBuf); return 0; } @@ -185,9 +220,13 @@ int MsgCenter_Send(const Msg_t *pstMsg) return RD_FAILURE; } - // MPSC安全写入(ringbuffer内部关中断+自旋锁) - int iRet = rd_RingbufferPutSafe(g_astModules[iDstIndex].pstRingBuf, - (const char *)pstMsg, sizeof(Msg_t)); + // 加锁保护ringbuffer(FreeRTOS: mutex, 裸机: 关中断+自旋锁) + MSG_CENTER_LOCK(&g_astModules[iDstIndex]); + + int iRet = rd_RingbufferPut(g_astModules[iDstIndex].pstRingBuf, + (const char *)pstMsg, sizeof(Msg_t)); + + MSG_CENTER_UNLOCK(&g_astModules[iDstIndex]); if (iRet > 0) { @@ -195,7 +234,7 @@ int MsgCenter_Send(const Msg_t *pstMsg) Rd_SemPost(&g_astModules[iDstIndex].stSem, __func__, NULL); return RD_SUCCESS; } - + log_w("warning queue full send ID:%x to %d", pstMsg->m_uiMsgID, pstMsg->m_uiDstModule); return RD_FAILURE; // 队列满 } @@ -276,12 +315,19 @@ int MsgCenter_SendFromISR(const char *pszDstName, uint32_t uiMsgID, const void * stMsg.m_uiDataLen = uiDataLen; } - // MPSC安全写入(ringbuffer内部关中断+自旋锁) - int iRet = rd_RingbufferPutSafe(g_astModules[iDstIndex].pstRingBuf, - (const char *)&stMsg, sizeof(Msg_t)); + // ISR安全:始终使用关中断+自旋锁(不依赖FreeRTOS mutex) + volatile uint8_t *pLock = &g_astModules[iDstIndex].bSpinLock; + uint32_t primask = rd_irq_save(); + while (ATOMIC_EXCH(pLock, 1, uint8_t, ATOMIC_ORDER_ACQUIRE)) {} + + int iRet = rd_RingbufferPut(g_astModules[iDstIndex].pstRingBuf, + (const char *)&stMsg, sizeof(Msg_t)); + + ATOMIC_STORE(pLock, 0, uint8_t, ATOMIC_ORDER_RELEASE); + rd_irq_restore(primask); + if (iRet > 0) { - // 信号量在CMSIS-RTOS2中自动识别ISR上下文 Rd_SemPost(&g_astModules[iDstIndex].stSem, __func__, NULL); return RD_SUCCESS; } @@ -299,9 +345,14 @@ int MsgCenter_SendFromISR(const char *pszDstName, uint32_t uiMsgID, const void * *****************************************************************************/ static int MsgCenter_ProcessOne(int iIndex, Msg_t *pstMsg) { - int iRet = rd_RingbufferGet(g_astModules[iIndex].pstRingBuf, + // 加锁保护ringbuffer + MSG_CENTER_LOCK(&g_astModules[iIndex]); + + int iRet = rd_RingbufferGet(g_astModules[iIndex].pstRingBuf, (char *)pstMsg, sizeof(Msg_t)); + MSG_CENTER_UNLOCK(&g_astModules[iIndex]); + return (iRet > 0) ? 1 : 0; } diff --git a/library/ringbuffer/include/ringbuffer.h b/library/ringbuffer/include/ringbuffer.h index 612b496..813bdba 100644 --- a/library/ringbuffer/include/ringbuffer.h +++ b/library/ringbuffer/include/ringbuffer.h @@ -74,7 +74,6 @@ struct TRingBuffer unsigned short m_sReadIndex; unsigned char m_bWriteMirror; unsigned short m_sWriteIndex; - unsigned char m_bWriteLock; /* 写入自旋锁(关中断+原子CAS,ISR安全) */ int m_iBufsize; char *m_pcBufPtr; }; @@ -121,7 +120,6 @@ extern void rd_RingbufferInit(rd_ringbuf_t *rb, char *pool, int size); extern int rd_RingbufferPeak(rd_ringbuf_t *rb, char **ptr); extern int rd_RingbufferPeekLinear(rd_ringbuf_t *rb, char *out, int length); extern int rd_RingbufferPut(rd_ringbuf_t *rb, const char *ptr, int length); -extern int rd_RingbufferPutSafe(rd_ringbuf_t *rb, const char *ptr, int length); extern int rd_RingbufferPutchar(rd_ringbuf_t *rb, const char ch); extern int rd_RingbufferPutcharForce(rd_ringbuf_t *rb, const char ch); extern int rd_RingbufferPutForce(rd_ringbuf_t *rb, const char *ptr, int length); diff --git a/library/ringbuffer/ringbuffer.c b/library/ringbuffer/ringbuffer.c index d4352d9..3139447 100644 --- a/library/ringbuffer/ringbuffer.c +++ b/library/ringbuffer/ringbuffer.c @@ -69,7 +69,6 @@ void rd_RingbufferInit(rd_ringbuf_t *rb, char *pool, int size) ATOMIC_STORE(&rb->m_sWriteIndex, 0, unsigned short, ATOMIC_ORDER_RELAXED); ATOMIC_STORE(&rb->m_bReadMirror, 0, unsigned char, ATOMIC_ORDER_RELAXED); ATOMIC_STORE(&rb->m_bWriteMirror, 0, unsigned char, ATOMIC_ORDER_RELAXED); - ATOMIC_STORE(&rb->m_bWriteLock, 0, unsigned char, ATOMIC_ORDER_RELAXED); rb->m_pcBufPtr = pool; rb->m_iBufsize = RD_ALIGN_DOWN(size, RD_ALIGN_SIZE); @@ -79,7 +78,6 @@ void rd_RingbufferInit(rd_ringbuf_t *rb, char *pool, int size) int rd_RingbufferSpaceLen(rd_ringbuf_t *rb) { if (!rb) return -1; - int iRet = 0; unsigned short wi = ATOMIC_LOAD(&rb->m_sWriteIndex, unsigned short, ATOMIC_ORDER_ACQUIRE); unsigned short ri = ATOMIC_LOAD(&rb->m_sReadIndex, unsigned short, ATOMIC_ORDER_RELAXED); @@ -87,19 +85,9 @@ int rd_RingbufferSpaceLen(rd_ringbuf_t *rb) unsigned char rb_m = ATOMIC_LOAD(&rb->m_bReadMirror, unsigned char, ATOMIC_ORDER_RELAXED); if (wb == rb_m) { - iRet = rb->m_iBufsize - (wi - ri); + return rb->m_iBufsize - (wi - ri); } else { - iRet = ri - wi; - } - - if (iRet >= 0) - { - return iRet; - } - else - { - log_e("wb = %d, rb_m = %d, wi = %d, ri = %d", wb, rb_m, wi, ri); - return 0; + return ri - wi; } } @@ -156,30 +144,6 @@ int rd_RingbufferPut(rd_ringbuf_t *rb, const char *ptr, int length) return length; } -/** - * @brief MPSC-safe Put: 关中断+原子自旋锁保护写入,可安全用于 ISR 和多线程并发写入。 - */ -int rd_RingbufferPutSafe(rd_ringbuf_t *rb, const char *ptr, int length) -{ - if (!rb || !ptr || length <= 0) return RD_INVALUE; - - /* 关中断(若已在 ISR 中则 primask=1,末尾不恢复开中断) */ - uint32_t primask = rd_irq_save(); - - /* 原子自旋锁 */ - while (ATOMIC_EXCH(&rb->m_bWriteLock, 1, unsigned char, ATOMIC_ORDER_ACQUIRE)) - { - /* 等待锁释放 */ - } - - int iRet = rd_RingbufferPut(rb, ptr, length); - - ATOMIC_STORE(&rb->m_bWriteLock, 0, unsigned char, ATOMIC_ORDER_RELEASE); - rd_irq_restore(primask); - - return iRet; -} - /** * @brief Put a block of data into the ring buffer. If the capacity of ring buffer is insufficient, it will overwrite the existing data in the ring buffer. * @@ -449,7 +413,6 @@ void rd_RingbufferReset(rd_ringbuf_t *rb) ATOMIC_STORE(&rb->m_sWriteIndex, 0, unsigned short, ATOMIC_ORDER_RELAXED); ATOMIC_STORE(&rb->m_bReadMirror, 0, unsigned char, ATOMIC_ORDER_RELAXED); ATOMIC_STORE(&rb->m_bWriteMirror, 0, unsigned char, ATOMIC_ORDER_RELAXED); - ATOMIC_STORE(&rb->m_bWriteLock, 0, unsigned char, ATOMIC_ORDER_RELAXED); } /**