slyant 5 lat temu
rodzic
commit
1b4703e3b2

+ 20 - 0
SConscript

@@ -0,0 +1,20 @@
+from building import *
+Import('rtconfig')
+
+src   = []
+cwd   = GetCurrentDir()
+
+# add src files.
+if GetDepend('PKG_USING_TASK_MSG_BUS'):
+    src += Glob('src/task_msg_bus.c')
+
+if GetDepend('PKG_USING_TASK_MSG_BUS_SAMPLE'):
+    src += Glob('examples/task_msg_bus_sample.c')
+
+# add include path.
+path  = [cwd + '/inc']
+
+# add src and include to group.
+group = DefineGroup('task_msg_bus', src, depend = ['PKG_USING_TASK_MSG_BUS'], CPPPATH = path)
+
+Return('group')

+ 77 - 0
examples/task_msg_bus_sample.c

@@ -0,0 +1,77 @@
+#include <board.h>
+#include "task_msg_name_user_def.h" //这行放到#include "task_msg_bus.h"前面先定义
+#include "task_msg_bus.h"
+
+#define DBG_TAG "task.msg.bus.sample"
+#define DBG_LVL DBG_LOG
+#include <rtdbg.h>
+
+static void msg_wait_thread_entry(void *params)
+{
+    struct task_msg_args args;
+    while(1)
+    {
+        /*
+        //测试 task_msg_wait_until
+        task_msg_wait_until(TASK_MSG_NET_REDAY, RT_WAITING_FOREVER, &args);
+        LOG_D("task_msg_wait_until => args.msg_name:%d, args.msg_args_json:%s", args.msg_name, args.msg_args_json);
+        rt_free(args.msg_args_json);    
+        */
+
+        //测试 task_msg_wait_any
+        const enum task_msg_name name_list[4] = {TASK_MSG_OS_REDAY, TASK_MSG_NET_REDAY, TASK_MSG_NET_REDAY3, TASK_MSG_NET_REDAY5};
+        task_msg_wait_any(name_list, sizeof(name_list)/sizeof(enum task_msg_name), RT_WAITING_FOREVER, &args);
+        LOG_D("task_msg_wait_any => args.msg_name:%d, args.msg_args_json:%s", args.msg_name, args.msg_args_json);
+        if(args.msg_args_json)
+            rt_free(args.msg_args_json);
+    }    
+}
+
+static void msg_publish_thread_entry(void *params)
+{
+    static int i = 0;
+    while (1)
+    {
+        if(i % 3 == 0)
+        {
+            task_msg_publish(TASK_MSG_OS_REDAY, "{\"os_reday\":true}");
+        }
+        else if(i % 3 == 1)
+        {
+            task_msg_publish(TASK_MSG_NET_REDAY3, "{\"net_reday\":true,\"ip\":\"10.0.0.20\"}");
+        }
+        else
+        {
+            task_msg_publish(TASK_MSG_NET_REDAY5, "{\"net_reday5\":true,\"ip\":\"192.168.0.50\"}");
+        }   
+        
+        rt_thread_mdelay(1000);
+        i++;
+    }
+}
+
+static void net_reday_callback(task_msg_args_t args)
+{
+    LOG_D("net_reday_callback => args->msg_name:%d, args->msg_args_json:%s", args->msg_name, args->msg_args_json);
+}
+
+static void os_reday_callback(task_msg_args_t args)
+{
+    LOG_D("os_reday_callback => args->msg_name:%d, args->msg_args_json:%s", args->msg_name, args->msg_args_json);
+}
+
+static int task_msg_bus_sample(void)
+{
+    //初始化消息总线(线程栈大小, 优先级, 时间片)
+    task_msg_bus_init(1024, 25, 10);
+    //订阅消息
+    task_msg_subscribe(TASK_MSG_NET_REDAY5, net_reday_callback);
+    task_msg_subscribe(TASK_MSG_OS_REDAY, os_reday_callback);
+    //创建一个等待消息的线程
+    rt_thread_t t_wait = rt_thread_create("msg_wait", msg_wait_thread_entry, RT_NULL, 512, 20, 10);
+    rt_thread_startup(t_wait);
+    //创建一个发布消息的线程
+    rt_thread_t t_publish = rt_thread_create("msg_pub", msg_publish_thread_entry, RT_NULL, 512, 15, 10);
+    rt_thread_startup(t_publish);
+}
+INIT_APP_EXPORT(task_msg_bus_sample);

+ 15 - 0
examples/task_msg_name_user_def.h

@@ -0,0 +1,15 @@
+#ifndef TASK_MSG_NAME_USER_DEF_H_
+#define TASK_MSG_NAME_USER_DEF_H_
+
+#define TASK_MSG_NAME_USER_DEF
+enum task_msg_name{
+    TASK_MSG_OS_REDAY = 0,
+    TASK_MSG_NET_REDAY,
+    TASK_MSG_NET_REDAY1,
+    TASK_MSG_NET_REDAY2,
+    TASK_MSG_NET_REDAY3,
+    TASK_MSG_NET_REDAY4,
+    TASK_MSG_NET_REDAY5,
+    TASK_MSG_COUNT
+};
+#endif

+ 66 - 0
inc/task_msg_bus.h

@@ -0,0 +1,66 @@
+/*
+ * Copyright (c) 2006-2020, RT-Thread Development Team
+ *
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * Change Logs:
+ * Date           Author       Notes
+ * 2020-03-22     sly_ant      the first version
+ */
+#ifndef TASK_MSG_BUS_H_
+#define TASK_MSG_BUS_H_
+
+#include <rtthread.h>
+#include <rtdevice.h>
+
+#include "task_msg_name_def.h"
+
+struct task_msg_args
+{
+    enum task_msg_name msg_name;
+    char *msg_args_json;
+};
+typedef struct task_msg_args *task_msg_args_t;
+
+struct task_msg_args_node
+{
+    task_msg_args_t args;
+    rt_slist_t slist;
+};
+typedef struct task_msg_args_node *task_msg_args_node_t;
+
+struct task_msg_callback_node
+{
+    void(*callback)(const task_msg_args_t msg_args);
+    rt_slist_t slist;
+};
+typedef struct task_msg_callback_node *task_msg_callback_node_t;
+
+struct task_msg_wait_node
+{
+    enum task_msg_name msg_name;
+    struct rt_semaphore msg_sem;
+    char *msg_args_json;
+    rt_slist_t slist;
+};
+typedef struct task_msg_wait_node *task_msg_wait_node_t;
+
+struct task_msg_wait_any_node
+{
+    enum task_msg_name *msg_name_list;
+    rt_uint8_t msg_name_list_len;
+    struct rt_semaphore msg_sem;
+    enum task_msg_name msg_name;
+    char *msg_args_json;
+    rt_slist_t slist;
+};
+typedef struct task_msg_wait_any_node *task_msg_wait_any_node_t;
+
+rt_err_t task_msg_bus_init(rt_uint32_t stack_size, rt_uint8_t  priority, rt_uint32_t tick);
+rt_err_t task_msg_subscribe(enum task_msg_name msg_name, void(*callback)(task_msg_args_t msg_args));
+rt_err_t task_msg_unsubscribe(enum task_msg_name msg_name, void(*callback)(task_msg_args_t msg_args));
+rt_err_t task_msg_publish(enum task_msg_name msg_name, const char *args_json);
+rt_err_t task_msg_wait_until(enum task_msg_name msg_name, rt_uint32_t timeout, task_msg_args_t out_args);
+rt_err_t task_msg_wait_any(const enum task_msg_name *msg_name_list, rt_uint8_t msg_name_list_len, rt_uint32_t timeout, task_msg_args_t out_args);
+
+#endif /* TASK_MSG_BUS_H_ */

+ 20 - 0
inc/task_msg_name_def.h

@@ -0,0 +1,20 @@
+/*
+ * Copyright (c) 2006-2020, RT-Thread Development Team
+ *
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * Change Logs:
+ * Date           Author       Notes
+ * 2020-03-22     sly_ant      the first version
+ */
+#ifndef TASK_MSG_NAME_DEF_H_
+#define TASK_MSG_NAME_DEF_H_
+
+#ifndef TASK_MSG_NAME_USER_DEF
+enum task_msg_name{
+    TASK_MSG_OS_REDAY = 0,
+    TASK_MSG_COUNT
+};
+#endif
+
+#endif /* TASK_MSG_NAME_DEF_H_ */

+ 311 - 0
src/task_msg_bus.c

@@ -0,0 +1,311 @@
+/*
+ * Copyright (c) 2006-2020, RT-Thread Development Team
+ *
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * Change Logs:
+ * Date           Author       Notes
+ * 2020-03-22     sly_ant      the first version
+ */
+
+#include "task_msg_bus.h"
+
+#define DBG_TAG "task.msg.bus"
+#define DBG_LVL DBG_LOG
+#include <rtdbg.h>
+
+static rt_bool_t task_msg_bus_init_tag = RT_FALSE;
+static struct rt_semaphore msg_sem;
+static struct rt_mutex msg_lock;
+static struct rt_mutex cb_lock;
+static struct rt_mutex wt_lock;
+static struct rt_mutex wta_lock;
+static rt_slist_t callback_slist_array[TASK_MSG_COUNT];
+static rt_slist_t msg_slist = RT_SLIST_OBJECT_INIT(msg_slist);
+static rt_slist_t wait_slist = RT_SLIST_OBJECT_INIT(wait_slist);
+static rt_slist_t wait_any_slist = RT_SLIST_OBJECT_INIT(wait_any_slist);
+
+rt_err_t task_msg_wait_any(const enum task_msg_name *msg_name_list, rt_uint8_t msg_name_list_len, rt_uint32_t timeout, task_msg_args_t out_args)
+{
+    if(task_msg_bus_init_tag==RT_FALSE) return -RT_EINVAL;
+
+    char name[RT_NAME_MAX];
+    task_msg_wait_any_node_t node = rt_calloc(1, sizeof(struct task_msg_wait_any_node));
+    if(node == RT_NULL)
+    {
+        LOG_E("there is no memory available!");
+        return RT_ENOMEM;
+    }
+    rt_uint32_t count = rt_slist_len(&wait_any_slist);
+    rt_snprintf(name, RT_NAME_MAX, "wta_%d", count);
+    rt_sem_init(&(node->msg_sem), name, 0, RT_IPC_FLAG_FIFO);
+    node->msg_name_list = (enum task_msg_name *)msg_name_list;
+    node->msg_name_list_len = msg_name_list_len;
+    node->msg_args_json = RT_NULL;
+    rt_slist_init(&(node->slist));
+    rt_mutex_take(&wta_lock, RT_WAITING_FOREVER);
+    rt_slist_append(&wait_any_slist, &(node->slist));
+    rt_mutex_release(&wta_lock);
+
+    rt_err_t rst = rt_sem_take(&(node->msg_sem), timeout);
+    if(rst==RT_EOK && out_args)
+    {
+        out_args->msg_name = node->msg_name;
+        out_args->msg_args_json = RT_NULL;
+        if(node->msg_args_json) out_args->msg_args_json = rt_strdup(node->msg_args_json);
+    }
+    rt_mutex_take(&wta_lock, RT_WAITING_FOREVER);
+    rt_slist_remove(&wait_any_slist, &(node->slist));
+    rt_mutex_release(&wta_lock);
+    rt_sem_detach(&(node->msg_sem));
+    if(node->msg_args_json)
+        rt_free(node->msg_args_json);
+    rt_free(node);
+
+    return rst;
+}
+
+rt_err_t task_msg_wait_until(enum task_msg_name msg_name, rt_uint32_t timeout, task_msg_args_t out_args)
+{
+    if(task_msg_bus_init_tag==RT_FALSE) return -RT_EINVAL;
+
+    char name[RT_NAME_MAX];
+    task_msg_wait_node_t node = rt_calloc(1, sizeof(struct task_msg_wait_node));
+    if(node == RT_NULL)
+    {
+        LOG_E("there is no memory available!");
+        return RT_ENOMEM;
+    }
+    rt_uint32_t count = rt_slist_len(&wait_slist);
+    rt_snprintf(name, RT_NAME_MAX, "wt_%d", count);
+    rt_sem_init(&(node->msg_sem), name, 0, RT_IPC_FLAG_FIFO);
+    node->msg_name = msg_name;
+    node->msg_args_json = RT_NULL;
+    rt_slist_init(&(node->slist));
+    rt_mutex_take(&wt_lock, RT_WAITING_FOREVER);
+    rt_slist_append(&wait_slist, &(node->slist));
+    rt_mutex_release(&wt_lock);
+
+    rt_err_t rst = rt_sem_take(&(node->msg_sem), timeout);
+    if(rst==RT_EOK && out_args)
+    {
+        out_args->msg_name = msg_name;
+        out_args->msg_args_json = RT_NULL;
+        if(node->msg_args_json) out_args->msg_args_json = rt_strdup(node->msg_args_json);
+    }
+    rt_mutex_take(&wt_lock, RT_WAITING_FOREVER);
+    rt_slist_remove(&wait_slist, &(node->slist));
+    rt_mutex_release(&wt_lock);
+    rt_sem_detach(&(node->msg_sem));
+    if(node->msg_args_json)
+        rt_free(node->msg_args_json);
+    rt_free(node);
+
+    return rst;
+}
+
+rt_err_t task_msg_subscribe(enum task_msg_name msg_name, void(*callback)(task_msg_args_t msg_args))
+{
+    if(task_msg_bus_init_tag==RT_FALSE) return -RT_EINVAL;
+
+    task_msg_callback_node_t callback_node = rt_calloc(1, sizeof(struct task_msg_callback_node));
+    if(callback_node == RT_NULL)
+    {
+        LOG_E("there is no memory available!");
+        return RT_ENOMEM;
+    }
+    callback_node->callback = callback;
+    rt_slist_init(&(callback_node->slist));
+
+    rt_mutex_take(&cb_lock, RT_WAITING_FOREVER);
+    rt_bool_t find_tag = RT_FALSE;
+    task_msg_callback_node_t node;
+    rt_slist_for_each_entry(node, &callback_slist_array[msg_name], slist)
+    {
+        if(node->callback == callback)
+        {
+            find_tag = RT_TRUE;
+            break;
+        }
+    }
+    if(find_tag)
+    {
+        LOG_W("this task msg callback with msg_name[%d] is exist!", msg_name);
+    }
+    else
+    {
+        rt_slist_append(&callback_slist_array[msg_name], &(callback_node->slist));
+    }
+    rt_mutex_release(&cb_lock);
+
+    return RT_EOK;
+}
+
+rt_err_t task_msg_unsubscribe(enum task_msg_name msg_name, void(*callback)(task_msg_args_t msg_args))
+{
+    if(task_msg_bus_init_tag==RT_FALSE) return -RT_EINVAL;
+
+    task_msg_callback_node_t node;
+    rt_mutex_take(&cb_lock, RT_WAITING_FOREVER);
+    rt_slist_for_each_entry(node, &callback_slist_array[msg_name], slist)
+    {
+        if(node->callback == callback)
+        {
+            rt_slist_remove(&callback_slist_array[msg_name], &(node->slist));
+            rt_free(node);
+            break;
+        }
+    }
+    rt_mutex_release(&cb_lock);
+
+    return RT_EOK;
+}
+
+rt_err_t task_msg_publish(enum task_msg_name msg_name, const char *args_json)
+{
+    if(task_msg_bus_init_tag==RT_FALSE) return -RT_EINVAL;
+
+    task_msg_args_node_t node = rt_calloc(1, sizeof(struct task_msg_args_node));
+    if(node == RT_NULL)
+    {
+        LOG_E("task msg publish failed! args_node create failed!");
+        return -RT_ENOMEM;
+    }
+
+    task_msg_args_t msg_args = rt_calloc(1, sizeof(struct task_msg_args));
+    if(msg_args == RT_NULL)
+    {
+        LOG_E("task msg publish failed! msg_args create failed!");
+        return -RT_ENOMEM;
+    }
+
+    msg_args->msg_name = msg_name;
+    msg_args->msg_args_json = RT_NULL;
+    if(args_json) msg_args->msg_args_json = rt_strdup(args_json);
+    node->args = msg_args;
+    rt_slist_init(&(node->slist));
+
+    rt_mutex_take(&msg_lock, RT_WAITING_FOREVER);
+    rt_slist_append(&msg_slist, &(node->slist));
+    rt_mutex_release(&msg_lock);
+
+    rt_sem_release(&msg_sem);
+
+    return RT_EOK;
+}
+
+static void task_msg_callback_init(void)
+{
+    rt_mutex_take(&cb_lock, RT_WAITING_FOREVER);
+    for(int i=0; i<TASK_MSG_COUNT; i++)
+    {
+        callback_slist_array[i].next = RT_NULL;
+    }
+    rt_mutex_release(&cb_lock);
+}
+
+static void task_msg_bus_thread_entry(void *params)
+{
+    while(1)
+    {
+        if(rt_sem_take(&msg_sem, RT_WAITING_FOREVER) == RT_EOK)
+        {
+            task_msg_args_node_t msg_args_node;
+            task_msg_callback_node_t msg_callback_node;
+            task_msg_wait_node_t msg_wait_node;
+            task_msg_wait_any_node_t msg_wait_any_node;
+            while(rt_slist_len(&msg_slist) > 0)
+            {
+                //get msg
+                msg_args_node = rt_slist_first_entry(&msg_slist, struct task_msg_args_node, slist);
+                //msg callback
+                rt_mutex_take(&cb_lock, RT_WAITING_FOREVER);
+                rt_slist_for_each_entry(msg_callback_node, &callback_slist_array[msg_args_node->args->msg_name], slist)
+                {
+                    if(msg_callback_node->callback) msg_callback_node->callback(msg_args_node->args);
+                }
+                rt_mutex_release(&cb_lock);
+                //check wait msg until
+                rt_mutex_take(&wt_lock, RT_WAITING_FOREVER);
+                rt_slist_for_each_entry(msg_wait_node, &wait_slist, slist)
+                {
+                    if(msg_wait_node->msg_name==msg_args_node->args->msg_name)
+                    {
+                        if(msg_args_node->args->msg_args_json)
+                            msg_wait_node->msg_args_json = rt_strdup(msg_args_node->args->msg_args_json);
+                        rt_sem_release(&(msg_wait_node->msg_sem));
+                    }
+                }
+                rt_mutex_release(&wt_lock);
+                //check wait any msg in array until
+                rt_mutex_take(&wta_lock, RT_WAITING_FOREVER);
+                rt_slist_for_each_entry(msg_wait_any_node, &wait_any_slist, slist)
+                {
+                    for(int i=0; i<msg_wait_any_node->msg_name_list_len; i++)
+                    {
+                        if(msg_wait_any_node->msg_name_list[i]==msg_args_node->args->msg_name)
+                        {
+                            msg_wait_any_node->msg_name = msg_args_node->args->msg_name;
+                            if(msg_args_node->args->msg_args_json)
+                                msg_wait_any_node->msg_args_json = rt_strdup(msg_args_node->args->msg_args_json);
+                            rt_sem_release(&(msg_wait_any_node->msg_sem));
+                            break;
+                        }
+                    }
+                }
+                rt_mutex_release(&wta_lock);
+                //remove msg
+                rt_mutex_take(&msg_lock, RT_WAITING_FOREVER);
+                rt_slist_remove(&msg_slist, &(msg_args_node->slist));
+                rt_mutex_release(&msg_lock);
+                //free msg
+                if(msg_args_node->args)
+                {
+                    if(msg_args_node->args->msg_args_json)
+                    {
+                        rt_free(msg_args_node->args->msg_args_json);
+                    }
+                    rt_free(msg_args_node->args);
+                }
+                rt_free(msg_args_node);
+            }
+        }
+    }
+}
+
+rt_err_t task_msg_bus_init(rt_uint32_t stack_size, rt_uint8_t  priority, rt_uint32_t tick)
+{
+    if(task_msg_bus_init_tag==RT_FALSE)
+    {
+        rt_sem_init(&msg_sem, "msg_sem", 0, RT_IPC_FLAG_FIFO);
+        rt_mutex_init(&msg_lock, "msg_lock", RT_IPC_FLAG_FIFO);
+        rt_mutex_init(&cb_lock, "cb_lock", RT_IPC_FLAG_FIFO);
+        rt_mutex_init(&wt_lock, "wt_lock", RT_IPC_FLAG_FIFO);
+        rt_mutex_init(&wta_lock, "wta_lock", RT_IPC_FLAG_FIFO);
+        task_msg_callback_init();
+        task_msg_bus_init_tag = RT_TRUE;
+    }
+
+    rt_thread_t t = rt_thread_create("msg_bus",
+            task_msg_bus_thread_entry,
+            RT_NULL,
+            stack_size,
+            priority,
+            tick);
+    if(t==RT_NULL)
+    {
+        LOG_E("task msg bus initialize failed! msg_bus_thread create failed!");
+        return -RT_ENOMEM;
+    }
+
+    rt_err_t rst = rt_thread_startup(t);
+    if(rst == RT_EOK)
+    {
+        LOG_I("task msg bus initialize success!");
+    }
+    else
+    {
+        LOG_E("task msg bus initialize failed! msg_bus thread startup failed(%d)", rst);
+    }
+    return rst;
+}