Browse Source

ringbuffer实现自旋锁(H7下)消息中心增加安全版本

paint_robot_new-v1.5
Lizongdi 6 days ago
parent
commit
4166f31136
  1. 10
      RBcore/include/msg_center.h
  2. 70
      RBcore/msg_center.c
  3. 16
      library/common/include/common.h
  4. 2
      library/ringbuffer/include/ringbuffer.h
  5. 26
      library/ringbuffer/ringbuffer.c

10
RBcore/include/msg_center.h

@ -109,6 +109,16 @@ int MsgCenter_Send(const Msg_t *pstMsg);
*/
int MsgCenter_SendTo(const char *pszDstName, uint32_t uiMsgID, const void *pData, uint32_t uiDataLen);
/**
* @brief 发送消息给指定模块(ISR安全版本,无mutex,关中断+自旋锁保护)
* @param pszDstName 目标模块名称
* @param uiMsgID 消息ID
* @param pData 数据指针(可为NULL)
* @param uiDataLen 数据长度
* @return 0成功,非0失败
*/
int MsgCenter_SendFromISR(const char *pszDstName, uint32_t uiMsgID, const void *pData, uint32_t uiDataLen);
/**
* @brief 处理当前模块的消息(在模块线程中调用)
* @param uiModuleID 模块ID

70
RBcore/msg_center.c

@ -45,7 +45,6 @@
typedef struct {
MsgModule_t stModule;
rd_ringbuf_t *pstRingBuf; // 消息队列(ringbuffer)
rd_mutex_t stMutex; // 互斥锁(保护ringbuffer)
rd_sem_t stSem; // 信号量(阻塞等待)
} MsgModuleEntry_t;
@ -129,17 +128,9 @@ uint32_t MsgCenter_Register(const char *pszName, MsgHandler_t pfHandler)
return 0;
}
// 初始化互斥锁
if (Rd_MutexInit(&g_astModules[uiIndex].stMutex) != RD_SUCCESS)
{
rd_RingbufferDestroy(g_astModules[uiIndex].pstRingBuf);
return 0;
}
// 初始化信号量(初始值为0,表示无消息)
if (Rd_SemInit(&g_astModules[uiIndex].stSem, 0, 0) != RD_SUCCESS)
{
Rd_MutexDestroy(&g_astModules[uiIndex].stMutex);
rd_RingbufferDestroy(g_astModules[uiIndex].pstRingBuf);
return 0;
}
@ -194,15 +185,10 @@ int MsgCenter_Send(const Msg_t *pstMsg)
return RD_FAILURE;
}
// 加锁保护ringbuffer
Rd_MutexLock(&g_astModules[iDstIndex].stMutex, __func__, NULL);
// 写入ringbuffer
int iRet = rd_RingbufferPut(g_astModules[iDstIndex].pstRingBuf,
// MPSC安全写入(ringbuffer内部关中断+自旋锁)
int iRet = rd_RingbufferPutSafe(g_astModules[iDstIndex].pstRingBuf,
(const char *)pstMsg, sizeof(Msg_t));
Rd_MutexUnlock(&g_astModules[iDstIndex].stMutex, __func__, NULL);
if (iRet > 0)
{
// 发送信号量,通知有新消息
@ -256,6 +242,53 @@ int MsgCenter_SendTo(const char *pszDstName, uint32_t uiMsgID, const void *pData
return MsgCenter_Send(&stMsg);
}
int MsgCenter_SendFromISR(const char *pszDstName, uint32_t uiMsgID, const void *pData, uint32_t uiDataLen)
{
if (NULL == pszDstName)
{
return RD_NULL;
}
uint32_t uiDstModuleID = MsgCenter_FindModule(pszDstName);
if (0 == uiDstModuleID)
{
return RD_FAILURE;
}
int iDstIndex = FIND_MODULE_INDEX(uiDstModuleID);
if (iDstIndex < 0 || NULL == g_astModules[iDstIndex].pstRingBuf)
{
return RD_FAILURE;
}
Msg_t stMsg;
RD_MEMSET(&stMsg, 0, sizeof(Msg_t));
stMsg.m_uiDstModule = uiDstModuleID;
stMsg.m_uiMsgID = uiMsgID;
if (NULL != pData && uiDataLen > 0)
{
if (uiDataLen > MSG_CENTER_MAX_DATA_SIZE)
{
uiDataLen = MSG_CENTER_MAX_DATA_SIZE;
}
RD_MEMCPY(stMsg.m_aucData, pData, uiDataLen);
stMsg.m_uiDataLen = uiDataLen;
}
// MPSC安全写入(ringbuffer内部关中断+自旋锁)
int iRet = rd_RingbufferPutSafe(g_astModules[iDstIndex].pstRingBuf,
(const char *)&stMsg, sizeof(Msg_t));
if (iRet > 0)
{
// 信号量在CMSIS-RTOS2中自动识别ISR上下文
Rd_SemPost(&g_astModules[iDstIndex].stSem, __func__, NULL);
return RD_SUCCESS;
}
return RD_FAILURE;
}
/*****************************************************************************
函 数 名 : MsgCenter_ProcessOne
功能描述 : 从ringbuffer读取一条消息
@ -266,14 +299,9 @@ int MsgCenter_SendTo(const char *pszDstName, uint32_t uiMsgID, const void *pData
*****************************************************************************/
static int MsgCenter_ProcessOne(int iIndex, Msg_t *pstMsg)
{
// 加锁保护ringbuffer
Rd_MutexLock(&g_astModules[iIndex].stMutex, __func__, NULL);
int iRet = rd_RingbufferGet(g_astModules[iIndex].pstRingBuf,
(char *)pstMsg, sizeof(Msg_t));
Rd_MutexUnlock(&g_astModules[iIndex].stMutex, __func__, NULL);
return (iRet > 0) ? 1 : 0;
}

16
library/common/include/common.h

@ -41,6 +41,22 @@
#define ATOMIC_ORDER_RELEASE 0
#endif
/* 中断控制(CMSIS 标准接口,ISR 安全) */
#if defined(__arm__) || defined(__aarch64__)
#include "../../../Drivers/CMSIS/Include/cmsis_compiler.h"
static inline uint32_t rd_irq_save(void) {
uint32_t primask = __get_PRIMASK();
__disable_irq();
return primask;
}
static inline void rd_irq_restore(uint32_t primask) {
__set_PRIMASK(primask);
}
#else
static inline uint32_t rd_irq_save(void) { return 0; }
static inline void rd_irq_restore(uint32_t primask) { (void)primask; }
#endif
/* 跨平台弱符号宏定义 */
#if defined(_MSC_VER) || defined(WIN32)
/* Microsoft Visual C++ */

2
library/ringbuffer/include/ringbuffer.h

@ -74,6 +74,7 @@ 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;
};
@ -120,6 +121,7 @@ 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);

26
library/ringbuffer/ringbuffer.c

@ -69,6 +69,7 @@ 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);
@ -155,6 +156,30 @@ 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.
*
@ -424,6 +449,7 @@ 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);
}
/**

Loading…
Cancel
Save