system/zbus: Port the Zephyr zbus message bus to NuttX

Port of the Zephyr RTOS zbus (many-to-many message bus with typed
channels and decoupled observers), built entirely on native NuttX
primitives and preserving the original declarative API
(ZBUS_CHAN_DEFINE, ZBUS_LISTENER_DEFINE, ZBUS_SUBSCRIBER_DEFINE, ...).

Features: listeners (synchronous callbacks), subscribers (queue of
channel references), message subscribers (ordered message copies),
async listeners (callback on a dedicated task), runtime observers,
per-observation notification masks, observer enable/disable, message
validators, channel user data, publish statistics, lookup by
name/numeric id and channel/observer iteration.

Mapping to NuttX primitives:
- Channel/observer registration: link-time iterable sections
  (include/nuttx/iterable_sections.h); the observers of a channel are
  named after their position in the definition, so the linker sorts the
  notification order, and the declarative macros are built on
  nuttx/macro.h (CONCATENATE, FOREACH_ARG and the FOREACH_IDX_ARG added
  in a companion nuttx commit) rather than on a private macro engine.
  Notification masks live in .bss with their initial value preserved in
  ROM and applied on lazy init.
- Channel lock: sem_t (enable CONFIG_PRIORITY_INHERITANCE instead of
  the Zephyr priority-boost/HLP); timeouts are computed with the
  clock_timespec_* helpers from nuttx/clock.h.
- Subscriber queues: kernel message queues (file_mq_*) opened lazily
  via pthread_once, usable from any task; mq payload copying replaces
  the Zephyr net_buf machinery entirely.
- Async listeners: one task per listener (task_create, priority and
  stack size configurable) blocking on the listener queue; a task
  rather than a pthread so it outlives the first API caller.
- Timeouts: milliseconds with CLOCK_MONOTONIC deadlines
  (ZBUS_NO_WAIT/ZBUS_FOREVER).

Includes a runnable example (examples/zbus, CONFIG_EXAMPLES_ZBUS) and a
cmocka test suite (testing/zbus, CONFIG_TESTING_ZBUS) covering the full
API: 17/17 tests passing on linum-stm32h753bi hardware, including
multi-channel index grouping, mask semantics, runtime observer error
paths, notification order (the observers of a channel run in the order
they are listed, and an observation bound with ZBUS_CHAN_ADD_OBS() runs
after all of them), queue overflow/timeout semantics, async listener
bursts,
bit-exact float/double payload delivery across every observer type
(sensor-style messages with a float-math validator) and an
interrupt-driven publisher (kernel timer interrupt -> signal -> sampling
thread -> zbus_chan_pub, the recommended pattern for interrupt sources).

Requirements: FLAT build; CONFIG_MQ_MAXMSGSIZE >= pointer size +
CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE for message subscribers; board
linker script including <nuttx/linker/common-rom.ld> or the generic
CONFIG_ITERABLE_SECTIONS_LINKER_INSERT mode.

Not ported: multi-domain proxy agent (experimental upstream); publishing
from interrupt handlers (userspace library: hand the data to a thread).

Documentation lives in the nuttx repository
(Documentation/applications/system/zbus).

Assisted-by: Claude Code
Signed-off-by: Jorge Guzman <jorge.gzm@gmail.com>
This commit is contained in:
Jorge Guzman 2026-07-10 09:23:58 -03:00 • committed by Alan C. Assis
parent 11124141cb
commit ee20ddd0ba
19 changed files with 3615 additions and 0 deletions

View file

@ -0,0 +1,33 @@
# ##############################################################################
# apps/examples/zbus/CMakeLists.txt
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more contributor
# license agreements. See the NOTICE file distributed with this work for
# additional information regarding copyright ownership. The ASF licenses this
# file to you under the Apache License, Version 2.0 (the "License"); you may not
# use this file except in compliance with the License. You may obtain a copy of
# the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations under
# the License.
#
# ##############################################################################
if(CONFIG_EXAMPLES_ZBUS)
nuttx_add_application(
NAME
${CONFIG_EXAMPLES_ZBUS_PROGNAME}
SRCS
zbus_main.c
STACKSIZE
${CONFIG_EXAMPLES_ZBUS_STACKSIZE}
PRIORITY
${CONFIG_EXAMPLES_ZBUS_PRIORITY})
endif()

27
examples/zbus/Kconfig Normal file
View file

@ -0,0 +1,27 @@
#
# For a description of the syntax of this configuration file,
# see the file kconfig-language.txt in the NuttX tools repository.
#
config EXAMPLES_ZBUS
tristate "ZBus example"
default n
depends on ZBUS
---help---
Enable the zbus message bus example.
if EXAMPLES_ZBUS
config EXAMPLES_ZBUS_PROGNAME
string "Program name"
default "zbus"
config EXAMPLES_ZBUS_PRIORITY
int "ZBus example task priority"
default 100
config EXAMPLES_ZBUS_STACKSIZE
int "ZBus example stack size"
default DEFAULT_TASK_STACKSIZE
endif

25
examples/zbus/Make.defs Normal file
View file

@ -0,0 +1,25 @@
############################################################################
# apps/examples/zbus/Make.defs
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership. The
# ASF licenses this file to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance with the
# License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations
# under the License.
#
############################################################################
ifneq ($(CONFIG_EXAMPLES_ZBUS),)
CONFIGURED_APPS += $(APPDIR)/examples/zbus
endif

34
examples/zbus/Makefile Normal file
View file

@ -0,0 +1,34 @@
############################################################################
# apps/examples/zbus/Makefile
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership. The
# ASF licenses this file to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance with the
# License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations
# under the License.
#
############################################################################
include $(APPDIR)/Make.defs
# ZBus example built-in application info
PROGNAME = $(CONFIG_EXAMPLES_ZBUS_PROGNAME)
PRIORITY = $(CONFIG_EXAMPLES_ZBUS_PRIORITY)
STACKSIZE = $(CONFIG_EXAMPLES_ZBUS_STACKSIZE)
MODULE = $(CONFIG_EXAMPLES_ZBUS)
MAINSRC = zbus_main.c
include $(APPDIR)/Application.mk

160
examples/zbus/zbus_main.c Normal file
View file

@ -0,0 +1,160 @@
/****************************************************************************
* apps/examples/zbus/zbus_main.c
*
* SPDX-License-Identifier: Apache-2.0
*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership. The
* ASF licenses this file to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance with the
* License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
* License for the specific language governing permissions and limitations
* under the License.
*
****************************************************************************/
/****************************************************************************
* Included Files
****************************************************************************/
#include <nuttx/config.h>
#include <pthread.h>
#include <stdio.h>
#include <unistd.h>
#include <system/zbus.h>
/****************************************************************************
* Private Types
****************************************************************************/
struct acc_msg
{
int x;
int y;
int z;
};
/****************************************************************************
* Private Function Prototypes
****************************************************************************/
static void listener_callback(const struct zbus_channel *chan);
/****************************************************************************
* Channel and observer definitions
****************************************************************************/
ZBUS_LISTENER_DEFINE(acc_listener, listener_callback);
ZBUS_SUBSCRIBER_DEFINE(acc_subscriber, 4);
ZBUS_CHAN_DEFINE(acc_chan, /* Name */
struct acc_msg, /* Message type */
NULL, /* Validator */
NULL, /* User data */
ZBUS_OBSERVERS(acc_listener, /* Observers */
acc_subscriber),
ZBUS_MSG_INIT(.x = 0, .y = 0, .z = 0));
/****************************************************************************
* Private Functions
****************************************************************************/
static void listener_callback(const struct zbus_channel *chan)
{
const struct acc_msg *msg = zbus_chan_const_msg(chan);
printf("zbus: listener: x=%d y=%d z=%d\n", msg->x, msg->y, msg->z);
}
static void *subscriber_thread(void *arg)
{
const struct zbus_channel *chan;
struct acc_msg msg;
int i;
for (i = 0; i < 5; i++)
{
if (zbus_sub_wait(&acc_subscriber, &chan, 2000) != 0)
{
printf("zbus: subscriber: timeout!\n");
continue;
}
if (chan == &acc_chan)
{
zbus_chan_read(chan, &msg, 500);
printf("zbus: subscriber: x=%d y=%d z=%d\n",
msg.x, msg.y, msg.z);
}
}
return NULL;
}
/****************************************************************************
* Public Functions
****************************************************************************/
int main(int argc, char *argv[])
{
struct acc_msg msg;
pthread_t thread;
int ret;
int i;
printf("zbus: publishing 5 messages to acc_chan\n");
ret = pthread_create(&thread, NULL, subscriber_thread, NULL);
if (ret != 0)
{
printf("zbus: could not create subscriber thread: %d\n", ret);
return 1;
}
for (i = 1; i <= 5; i++)
{
msg.x = i;
msg.y = i * 10;
msg.z = i * 100;
ret = zbus_chan_pub(&acc_chan, &msg, 1000);
if (ret != 0)
{
printf("zbus: publish error: %d\n", ret);
}
/* Mask the listener notifications on the third message to
* demonstrate the notification mask API.
*/
if (i == 3)
{
zbus_obs_set_chan_notification_mask(&acc_listener, &acc_chan,
true);
printf("zbus: listener masked\n");
}
else if (i == 4)
{
zbus_obs_set_chan_notification_mask(&acc_listener, &acc_chan,
false);
printf("zbus: listener unmasked\n");
}
usleep(100 * 1000);
}
pthread_join(thread, NULL);
printf("zbus: done\n");
return 0;
}

719
include/system/zbus.h Normal file
View file

@ -0,0 +1,719 @@
/****************************************************************************
* apps/include/system/zbus.h
*
* SPDX-License-Identifier: Apache-2.0
*
* Copyright (c) 2022 Rodrigo Peixoto <rodrigopex@gmail.com>
* Copyright (c) 2026 NuttX port
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may
* not use this file except in compliance with the License. You may obtain
* a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
* License for the specific language governing permissions and limitations
* under the License.
*
****************************************************************************/
/* NuttX port of the Zephyr zbus message bus.
*
* Differences from the Zephyr original:
* - Timeouts are given in milliseconds (int32_t): ZBUS_NO_WAIT (0) and
* ZBUS_FOREVER (-1) replace K_NO_WAIT/K_FOREVER.
* - Subscribers and message subscribers use POSIX message queues opened
* lazily on first zbus API call (no k_msgq/k_fifo/net_buf).
* - Priority boost (HLP) is not implemented; enable NuttX native
* CONFIG_PRIORITY_INHERITANCE for equivalent protection.
* - Publishing from interrupt context is not supported.
* - Requires the board linker script to include the iterable section
* fragments <nuttx/linker/common-rom.ld> and common-ram.ld.
*/
#ifndef __APPS_INCLUDE_SYSTEM_ZBUS_H
#define __APPS_INCLUDE_SYSTEM_ZBUS_H
/****************************************************************************
* Included Files
****************************************************************************/
#include <nuttx/config.h>
#include <nuttx/fs/fs.h>
#include <nuttx/iterable_sections.h>
#include <assert.h>
#include <errno.h>
#include <semaphore.h>
#include <stdbool.h>
#include <stddef.h>
#include <stdint.h>
#include <string.h>
#include <time.h>
#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
# include <nuttx/list.h>
#endif
#include <sys/types.h>
#include <nuttx/macro.h>
#ifdef __cplusplus
#define _ZBUS_CPP_EXTERN extern
extern "C"
{
#else
#define _ZBUS_CPP_EXTERN
#endif
/****************************************************************************
* Pre-processor Definitions
****************************************************************************/
/* Timeout special values (milliseconds) */
#define ZBUS_NO_WAIT 0
#define ZBUS_FOREVER (-1)
/* Channel without a unique numeric identifier */
#define ZBUS_CHAN_ID_INVALID UINT32_MAX
#ifdef CONFIG_ZBUS_ASSERT_MOCK
# define _ZBUS_ASSERT(cond, msg) \
do \
{ \
if (!(cond)) \
{ \
return -EFAULT; \
} \
} \
while (0)
#else
# define _ZBUS_ASSERT(cond, msg) DEBUGASSERT(cond)
#endif
/****************************************************************************
* Public Types
****************************************************************************/
struct zbus_channel;
/* Mutable data associated with every channel */
struct zbus_channel_data
{
/* Boundaries of this channel's static observations inside the sorted
* zbus_channel_observation iterable section (computed on first use).
*/
int16_t observers_start_idx;
int16_t observers_end_idx;
/* Channel access semaphore */
sem_t sem;
#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
/* Runtime (dynamically added) observers */
struct list_node observers;
#endif
#ifdef CONFIG_ZBUS_CHANNEL_PUBLISH_STATS
struct timespec publish_timestamp;
uint32_t publish_count;
#endif
};
/* A channel: constant descriptor placed in ROM (iterable section) */
struct zbus_channel
{
#ifdef CONFIG_ZBUS_CHANNEL_NAME
const char *name;
#endif
#ifdef CONFIG_ZBUS_CHANNEL_ID
uint32_t id;
#endif
/* Shared message memory, its size, and optional user data/validator */
void *message;
size_t message_size;
void *user_data;
bool (*validator)(const void *msg, size_t msg_size);
struct zbus_channel_data *data;
};
/* Observer types */
enum zbus_observer_type
{
ZBUS_OBSERVER_LISTENER_TYPE = 0,
ZBUS_OBSERVER_SUBSCRIBER_TYPE,
ZBUS_OBSERVER_MSG_SUBSCRIBER_TYPE,
ZBUS_OBSERVER_ASYNC_LISTENER_TYPE,
};
/* Mutable data associated with every observer */
struct zbus_observer_data
{
bool enabled;
/* Notification queue (subscriber/msg subscriber/async listener), opened
* lazily with file_mq_open() so it is usable from any task, unlike
* per-task mqd_t descriptors. mq.f_inode == NULL means "not opened".
*/
struct file mq;
#ifdef CONFIG_ZBUS_ASYNC_LISTENER
/* Dedicated task running the async listener callback. A task (not a
* pthread) so it outlives the task that triggered the lazy init.
*/
pid_t pid;
#endif
};
/* An observer: constant descriptor placed in ROM (iterable section) */
struct zbus_observer
{
#ifdef CONFIG_ZBUS_OBSERVER_NAME
const char *name;
#endif
enum zbus_observer_type type;
/* Notification queue depth (subscriber types only) */
uint16_t queue_size;
struct zbus_observer_data *data;
/* Listener callback (listener type only) */
void (*callback)(const struct zbus_channel *chan);
#ifdef CONFIG_ZBUS_ASYNC_LISTENER
/* Async listener callback (async listener type only). Executed on the
* listener's dedicated task with a copy of the published message.
*/
void (*async_callback)(const struct zbus_channel *chan, const void *msg);
#endif
};
/* Link between one channel and one observer (ROM iterable section, sorted
* by name so that entries are grouped by channel and ordered by observer
* priority). The mutable notification mask lives in .bss and is pointed
* to from here; its initial value is preserved in ROM (mask_init) and
* applied by the one-time lazy initialization.
*/
struct zbus_channel_observation
{
const struct zbus_channel *chan;
const struct zbus_observer *obs;
bool *mask;
bool mask_init;
};
#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
/* Node linking a runtime observer to a channel */
struct zbus_observer_node
{
struct list_node node;
const struct zbus_observer *obs;
};
#endif
/****************************************************************************
* Definition macros
****************************************************************************/
#ifdef CONFIG_ZBUS_CHANNEL_NAME
# define ZBUS_CHANNEL_NAME_INIT(_name) .name = #_name,
#else
# define ZBUS_CHANNEL_NAME_INIT(_name)
#endif
#ifdef CONFIG_ZBUS_CHANNEL_ID
# define _ZBUS_CHANNEL_ID_INIT(_id) .id = _id,
#else
# define _ZBUS_CHANNEL_ID_INIT(_id)
#endif
#ifdef CONFIG_ZBUS_OBSERVER_NAME
# define ZBUS_OBSERVER_NAME_INIT(_name) .name = #_name,
#else
# define ZBUS_OBSERVER_NAME_INIT(_name)
#endif
#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
# define _ZBUS_RUNTIME_OBS_INIT(_name) \
.observers = LIST_INITIAL_VALUE(_zbus_chan_data_##_name.observers),
#else
# define _ZBUS_RUNTIME_OBS_INIT(_name)
#endif
#define _ZBUS_MESSAGE_NAME(_name) _zbus_message_##_name
/* Declare channels/observers defined in other files */
#define _ZBUS_OBS_EXTERN(_p, _name, _i) \
extern const struct zbus_observer _name;
#define _ZBUS_CHAN_EXTERN(_p, _name, _i) \
extern const struct zbus_channel _name;
#define ZBUS_OBS_DECLARE(...) \
FOREACH_ARG(_ZBUS_OBS_EXTERN, 0, __VA_ARGS__)
#define ZBUS_CHAN_DECLARE(...) \
FOREACH_ARG(_ZBUS_CHAN_EXTERN, 0, __VA_ARGS__)
/* Observer list helpers for ZBUS_CHAN_DEFINE */
#define ZBUS_OBSERVERS_EMPTY
#define ZBUS_OBSERVERS(...) __VA_ARGS__
/* Message initializer: ZBUS_MSG_INIT(.a = 1, .b = 2) -> {.a = 1, .b = 2} */
#define ZBUS_MSG_INIT(_val, ...) {_val, ##__VA_ARGS__}
/* One channel<->observer observation and its mask. The variable name
* embeds the channel and observer names, so the linker's SORT_BY_NAME()
* keeps the observations of a channel contiguous, and the index embedded
* in the name gives the notification order within the channel.
*/
/* FOREACH_ARG() hands the position as a plain literal, and the linker
* sorts the observations of a channel by name, so the index has to be
* zero padded to a fixed width for the sort to follow the declaration
* order beyond ten observers. Same table as the original zbus.
*/
#define _ZBUS_OBS_IDX_0 00
#define _ZBUS_OBS_IDX_1 01
#define _ZBUS_OBS_IDX_2 02
#define _ZBUS_OBS_IDX_3 03
#define _ZBUS_OBS_IDX_4 04
#define _ZBUS_OBS_IDX_5 05
#define _ZBUS_OBS_IDX_6 06
#define _ZBUS_OBS_IDX_7 07
#define _ZBUS_OBS_IDX_8 08
#define _ZBUS_OBS_IDX_9 09
#define _ZBUS_OBS_IDX_10 10
#define _ZBUS_OBS_IDX_11 11
#define _ZBUS_OBS_IDX_12 12
#define _ZBUS_OBS_IDX_13 13
#define _ZBUS_OBS_IDX_14 14
#define _ZBUS_OBS_IDX_15 15
#define _ZBUS_OBS_IDX_16 16
#define _ZBUS_OBS_IDX_17 17
#define _ZBUS_OBS_IDX_18 18
#define _ZBUS_OBS_IDX_19 19
#define _ZBUS_OBS_IDX_20 20
#define _ZBUS_OBS_IDX_21 21
#define _ZBUS_OBS_IDX_22 22
#define _ZBUS_OBS_IDX_23 23
#define _ZBUS_OBS_IDX_24 24
#define _ZBUS_OBS_IDX_25 25
#define _ZBUS_OBS_IDX_26 26
#define _ZBUS_OBS_IDX_27 27
#define _ZBUS_OBS_IDX_28 28
#define _ZBUS_OBS_IDX_29 29
#define _ZBUS_OBS_IDX_30 30
#define _ZBUS_OBS_IDX_31 31
#define _ZBUS_OBS_IDX(_idx) CONCATENATE(_ZBUS_OBS_IDX_, _idx)
#define _ZBUS_OBSERVATION_NAME_(_chan, _idx, _obs) \
_zbus_obn_##_chan##_##_idx##_##_obs
/* One level of indirection so that _idx is expanded before the paste,
* which is what lets the caller build it with CONCATENATE().
*/
#define _ZBUS_OBSERVATION_NAME(_chan, _idx, _obs) \
_ZBUS_OBSERVATION_NAME_(_chan, _idx, _obs)
/* The observation name reaches this macro already expanded, so the input
* section is named after the variable. It must not expand FOREACH_ARG():
* it is itself expanded from within one, and the preprocessor does not
* rescan a macro that is already being expanded.
*/
#define _ZBUS_OBSERVATION_DEFINE(_obn, _chan, _obs, _masked) \
static bool CONCATENATE(_obn, _mask) = _masked; \
const STRUCT_SECTION_ITERABLE(zbus_channel_observation, _obn) = \
{ \
.chan = &_chan, \
.obs = &_obs, \
.mask = &CONCATENATE(_obn, _mask), \
.mask_init = _masked, \
}
#define _ZBUS_CHAN_OBSERVATION(_chan, _obs, _idx) \
_ZBUS_OBSERVATION_DEFINE( \
_ZBUS_OBSERVATION_NAME(_chan, _ZBUS_OBS_IDX(_idx), _obs), \
_chan, _obs, false);
#define _ZBUS_CHAN_DEFINE(_name, _id, _type, _validator, _user_data) \
static struct zbus_channel_data _zbus_chan_data_##_name = \
{ \
.observers_start_idx = -1, \
.observers_end_idx = -1, \
.sem = SEM_INITIALIZER(1), \
_ZBUS_RUNTIME_OBS_INIT(_name) \
}; \
_ZBUS_CPP_EXTERN const STRUCT_SECTION_ITERABLE(zbus_channel, _name) = \
{ \
ZBUS_CHANNEL_NAME_INIT(_name) \
_ZBUS_CHANNEL_ID_INIT(_id) \
.message = &_ZBUS_MESSAGE_NAME(_name), \
.message_size = sizeof(_type), \
.user_data = _user_data, \
.validator = _validator, \
.data = &_zbus_chan_data_##_name, \
}
/* Define a channel.
*
* _name channel name (C identifier)
* _type message type (struct or union)
* _validator optional validator function or NULL
* _user_data optional user data pointer or NULL
* _observers ZBUS_OBSERVERS(obs1, obs2, ...) or ZBUS_OBSERVERS_EMPTY;
* the position in the list becomes the observation priority,
* so the observers are notified in the order listed
* _init_val message initial value, e.g. ZBUS_MSG_INIT(0)
*/
#define ZBUS_CHAN_DEFINE(_name, _type, _validator, _user_data, _observers, \
_init_val) \
static _type _ZBUS_MESSAGE_NAME(_name) = _init_val; \
_ZBUS_CHAN_DEFINE(_name, ZBUS_CHAN_ID_INVALID, _type, _validator, \
_user_data); \
ZBUS_OBS_DECLARE(_observers) \
FOREACH_ARG(_ZBUS_CHAN_OBSERVATION, _name, _observers)
/* Same as ZBUS_CHAN_DEFINE with a unique numeric channel identifier */
#define ZBUS_CHAN_DEFINE_WITH_ID(_name, _id, _type, _validator, _user_data, \
_observers, _init_val) \
static _type _ZBUS_MESSAGE_NAME(_name) = _init_val; \
_ZBUS_CHAN_DEFINE(_name, _id, _type, _validator, _user_data); \
ZBUS_OBS_DECLARE(_observers) \
FOREACH_ARG(_ZBUS_CHAN_OBSERVATION, _name, _observers)
/* Add a static observation to a channel defined elsewhere. _prio defines
* the notification order relative to other ADD_OBS observations of the
* same channel (use two-digit literals, e.g. 01, 02, ... so the linker
* name sort orders them correctly). ADD_OBS observations are notified
* after the ones listed in ZBUS_CHAN_DEFINE.
*/
/* Observations added out of line are notified after every observer listed
* in the channel definition, and _prio only orders them among themselves:
* the "zz" infix places them after the two digit indexes in the name
* sorted section, exactly as the original zbus does.
*/
#define ZBUS_CHAN_ADD_OBS_WITH_MASK(_chan, _obs, _masked, _prio) \
_ZBUS_OBSERVATION_DEFINE( \
_ZBUS_OBSERVATION_NAME(_chan, CONCATENATE(zz, _prio), _obs), \
_chan, _obs, _masked)
#define ZBUS_CHAN_ADD_OBS(_chan, _obs, _prio) \
ZBUS_CHAN_ADD_OBS_WITH_MASK(_chan, _obs, false, _prio)
/* Define a listener observer (synchronous callback) */
#define ZBUS_LISTENER_DEFINE_WITH_ENABLE(_name, _cb, _enable) \
static struct zbus_observer_data _zbus_obs_data_##_name = \
{ \
.enabled = _enable, \
}; \
_ZBUS_CPP_EXTERN const STRUCT_SECTION_ITERABLE(zbus_observer, _name) = \
{ \
ZBUS_OBSERVER_NAME_INIT(_name) \
.type = ZBUS_OBSERVER_LISTENER_TYPE, \
.queue_size = 0, \
.data = &_zbus_obs_data_##_name, \
.callback = (_cb), \
}
#define ZBUS_LISTENER_DEFINE(_name, _cb) \
ZBUS_LISTENER_DEFINE_WITH_ENABLE(_name, _cb, true)
/* Define a subscriber observer (receives channel references through a
* message queue of depth _queue_size; use zbus_sub_wait() to wait).
*/
#define ZBUS_SUBSCRIBER_DEFINE_WITH_ENABLE(_name, _queue_size, _enable) \
static struct zbus_observer_data _zbus_obs_data_##_name = \
{ \
.enabled = _enable, \
}; \
_ZBUS_CPP_EXTERN const STRUCT_SECTION_ITERABLE(zbus_observer, _name) = \
{ \
ZBUS_OBSERVER_NAME_INIT(_name) \
.type = ZBUS_OBSERVER_SUBSCRIBER_TYPE, \
.queue_size = _queue_size, \
.data = &_zbus_obs_data_##_name, \
.callback = NULL, \
}
#define ZBUS_SUBSCRIBER_DEFINE(_name, _queue_size) \
ZBUS_SUBSCRIBER_DEFINE_WITH_ENABLE(_name, _queue_size, true)
#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
/* Define a message subscriber observer (receives copies of the published
* messages through a message queue; use zbus_sub_wait_msg() to wait).
* Messages larger than CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE cannot be
* delivered to message subscribers.
*/
#define ZBUS_MSG_SUBSCRIBER_DEFINE_WITH_ENABLE(_name, _enable) \
static struct zbus_observer_data _zbus_obs_data_##_name = \
{ \
.enabled = _enable, \
}; \
_ZBUS_CPP_EXTERN const STRUCT_SECTION_ITERABLE(zbus_observer, _name) = \
{ \
ZBUS_OBSERVER_NAME_INIT(_name) \
.type = ZBUS_OBSERVER_MSG_SUBSCRIBER_TYPE, \
.queue_size = CONFIG_ZBUS_MSG_SUBSCRIBER_QUEUE_SIZE, \
.data = &_zbus_obs_data_##_name, \
.callback = NULL, \
}
#define ZBUS_MSG_SUBSCRIBER_DEFINE(_name) \
ZBUS_MSG_SUBSCRIBER_DEFINE_WITH_ENABLE(_name, true)
#endif /* CONFIG_ZBUS_MSG_SUBSCRIBER */
#ifdef CONFIG_ZBUS_ASYNC_LISTENER
/* Define an async listener observer. The callback executes on a
* dedicated task (not in the publisher context) and receives a copy of
* the published message. Messages larger than
* CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE cannot be delivered.
*/
#define ZBUS_ASYNC_LISTENER_DEFINE_WITH_ENABLE(_name, _cb, _enable) \
static struct zbus_observer_data _zbus_obs_data_##_name = \
{ \
.enabled = _enable, \
}; \
_ZBUS_CPP_EXTERN const STRUCT_SECTION_ITERABLE(zbus_observer, _name) = \
{ \
ZBUS_OBSERVER_NAME_INIT(_name) \
.type = ZBUS_OBSERVER_ASYNC_LISTENER_TYPE, \
.queue_size = CONFIG_ZBUS_MSG_SUBSCRIBER_QUEUE_SIZE, \
.data = &_zbus_obs_data_##_name, \
.callback = NULL, \
.async_callback = (_cb), \
}
#define ZBUS_ASYNC_LISTENER_DEFINE(_name, _cb) \
ZBUS_ASYNC_LISTENER_DEFINE_WITH_ENABLE(_name, _cb, true)
#endif /* CONFIG_ZBUS_ASYNC_LISTENER */
/****************************************************************************
* Public Function Prototypes
****************************************************************************/
/* Publish a message to a channel. Copies *msg into the channel and runs
* the dispatcher, notifying every observer. Returns 0 or -errno
* (-ENOMSG: validator rejected; -EBUSY/-EAGAIN: could not lock in time).
*/
int zbus_chan_pub(const struct zbus_channel *chan, const void *msg,
int32_t timeout_ms);
/* Read a channel message (copies the channel message into *msg) */
int zbus_chan_read(const struct zbus_channel *chan, void *msg,
int32_t timeout_ms);
/* Force the notification of a channel's observers without publishing */
int zbus_chan_notify(const struct zbus_channel *chan, int32_t timeout_ms);
/* Claim/finish a channel for direct access to zbus_chan_msg() */
int zbus_chan_claim(const struct zbus_channel *chan, int32_t timeout_ms);
int zbus_chan_finish(const struct zbus_channel *chan);
/* Wait for a notification (subscriber observers) */
int zbus_sub_wait(const struct zbus_observer *sub,
const struct zbus_channel **chan, int32_t timeout_ms);
#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
/* Wait for a message copy (message subscriber observers) */
int zbus_sub_wait_msg(const struct zbus_observer *sub,
const struct zbus_channel **chan, void *msg,
int32_t timeout_ms);
#endif
/* Enable/disable an observer */
int zbus_obs_set_enable(const struct zbus_observer *obs, bool enabled);
/* Mask/unmask the notifications from one channel to one observer */
int zbus_obs_set_chan_notification_mask(const struct zbus_observer *obs,
const struct zbus_channel *chan,
bool masked);
int zbus_obs_is_chan_notification_masked(const struct zbus_observer *obs,
const struct zbus_channel *chan,
bool *masked);
#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
/* Add/remove observers at runtime */
int zbus_chan_add_obs(const struct zbus_channel *chan,
const struct zbus_observer *obs, int32_t timeout_ms);
int zbus_chan_rm_obs(const struct zbus_channel *chan,
const struct zbus_observer *obs, int32_t timeout_ms);
#endif
#ifdef CONFIG_ZBUS_CHANNEL_ID
const struct zbus_channel *zbus_chan_from_id(uint32_t channel_id);
#endif
#ifdef CONFIG_ZBUS_CHANNEL_NAME
const struct zbus_channel *zbus_chan_from_name(const char *name);
#endif
/* Iteration over all channels/observers. The iterator function returns
* false to stop the iteration.
*/
bool zbus_iterate_over_channels(
bool (*iterator_func)(const struct zbus_channel *chan));
bool zbus_iterate_over_channels_with_user_data(
bool (*iterator_func)(const struct zbus_channel *chan, void *user_data),
void *user_data);
bool zbus_iterate_over_observers(
bool (*iterator_func)(const struct zbus_observer *obs));
bool zbus_iterate_over_observers_with_user_data(
bool (*iterator_func)(const struct zbus_observer *obs, void *user_data),
void *user_data);
/****************************************************************************
* Inline Functions
****************************************************************************/
#ifdef CONFIG_ZBUS_CHANNEL_NAME
static inline const char *zbus_chan_name(const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
return chan->name;
}
#endif
/* Direct access to the channel message. Only valid while the channel is
* locked (inside a listener callback or between claim/finish).
*/
static inline void *zbus_chan_msg(const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
return chan->message;
}
static inline const void *zbus_chan_const_msg(
const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
return chan->message;
}
static inline size_t zbus_chan_msg_size(const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
return chan->message_size;
}
static inline void *zbus_chan_user_data(const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
return chan->user_data;
}
static inline int zbus_obs_is_enabled(const struct zbus_observer *obs,
bool *enable)
{
_ZBUS_ASSERT(obs != NULL, "obs is required");
_ZBUS_ASSERT(enable != NULL, "enable is required");
*enable = obs->data->enabled;
return 0;
}
#ifdef CONFIG_ZBUS_OBSERVER_NAME
static inline const char *zbus_obs_name(const struct zbus_observer *obs)
{
DEBUGASSERT(obs != NULL);
return obs->name;
}
#endif
#ifdef CONFIG_ZBUS_CHANNEL_PUBLISH_STATS
/* Update the publish statistics (claim/finish workflow only; the channel
* must be locked).
*/
static inline void zbus_chan_pub_stats_update(
const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
clock_gettime(CLOCK_MONOTONIC, &chan->data->publish_timestamp);
chan->data->publish_count += 1;
}
static inline struct timespec zbus_chan_pub_stats_last_time(
const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
return chan->data->publish_timestamp;
}
static inline uint32_t zbus_chan_pub_stats_count(
const struct zbus_channel *chan)
{
DEBUGASSERT(chan != NULL);
return chan->data->publish_count;
}
#else
static inline void zbus_chan_pub_stats_update(
const struct zbus_channel *chan)
{
(void)chan;
}
#endif /* CONFIG_ZBUS_CHANNEL_PUBLISH_STATS */
#ifdef __cplusplus
}
#endif
#endif /* __APPS_INCLUDE_SYSTEM_ZBUS_H */

View file

@ -0,0 +1,31 @@
# ##############################################################################
# apps/system/zbus/CMakeLists.txt
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more contributor
# license agreements. See the NOTICE file distributed with this work for
# additional information regarding copyright ownership. The ASF licenses this
# file to you under the Apache License, Version 2.0 (the "License"); you may not
# use this file except in compliance with the License. You may obtain a copy of
# the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations under
# the License.
#
# ##############################################################################
if(CONFIG_ZBUS)
set(SRCS zbus.c zbus_iterable_sections.c)
if(CONFIG_ZBUS_RUNTIME_OBSERVERS)
list(APPEND SRCS zbus_runtime_observers.c)
endif()
target_sources(apps PRIVATE ${SRCS})
endif()

110
system/zbus/Kconfig Normal file
View file

@ -0,0 +1,110 @@
#
# For a description of the syntax of this configuration file,
# see the file kconfig-language.txt in the NuttX tools repository.
#
menuconfig ZBUS
bool "ZBus message bus library"
default n
depends on !DISABLE_MQUEUE
---help---
Enable the zbus message bus library (port of the Zephyr zbus).
Channels and observers are defined statically with the
ZBUS_CHAN_DEFINE/ZBUS_LISTENER_DEFINE/ZBUS_SUBSCRIBER_DEFINE
macros and collected in linker iterable sections. The board
linker script must include <nuttx/linker/common-rom.ld> (inside
.text) and <nuttx/linker/common-ram.ld> (inside .data).
For protection against priority inversion during the
notification process, enable CONFIG_PRIORITY_INHERITANCE.
if ZBUS
config ZBUS_CHANNEL_NAME
bool "Channel name field"
default n
---help---
Store the channel name string and enable zbus_chan_name() and
zbus_chan_from_name().
config ZBUS_CHANNEL_ID
bool "Channel identifier field"
default n
---help---
Store a unique numeric channel identifier and enable
zbus_chan_from_id(). Use ZBUS_CHAN_DEFINE_WITH_ID.
config ZBUS_OBSERVER_NAME
bool "Observer name field"
default n
---help---
Store the observer name string and enable zbus_obs_name().
config ZBUS_CHANNEL_PUBLISH_STATS
bool "Channel publishing statistics (timestamp and count)"
default n
config ZBUS_MSG_SUBSCRIBER
bool "Message subscribers (receive message copies in sequence)"
default n
---help---
Enable ZBUS_MSG_SUBSCRIBER_DEFINE and zbus_sub_wait_msg().
Message subscribers receive a copy of every published message
through a POSIX message queue.
if ZBUS_MSG_SUBSCRIBER
config ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE
int "Size of the biggest message used with zbus (bytes)"
default 64
---help---
Messages larger than this cannot be delivered to message
subscribers. Defines the message queue slot size.
NOTE: CONFIG_MQ_MAXMSGSIZE must be at least this value plus
the size of a pointer, otherwise the message subscriber
queues fail to open with -EINVAL.
config ZBUS_MSG_SUBSCRIBER_QUEUE_SIZE
int "Message subscriber queue depth"
default 4
endif # ZBUS_MSG_SUBSCRIBER
config ZBUS_ASYNC_LISTENER
bool "Async listeners"
default n
depends on ZBUS_MSG_SUBSCRIBER
---help---
Async listeners execute their callback on a dedicated task
(one per async listener, spawned on first use) with a copy of
the published message, instead of running synchronously in the
publisher context. Enable with ZBUS_ASYNC_LISTENER_DEFINE.
if ZBUS_ASYNC_LISTENER
config ZBUS_ASYNC_LISTENER_PRIORITY
int "Async listener task priority"
default 100
config ZBUS_ASYNC_LISTENER_STACKSIZE
int "Async listener task stack size"
default DEFAULT_TASK_STACKSIZE
endif # ZBUS_ASYNC_LISTENER
config ZBUS_RUNTIME_OBSERVERS
bool "Runtime observers support"
default n
---help---
Enable zbus_chan_add_obs()/zbus_chan_rm_obs(). Observer nodes
are allocated from the heap.
config ZBUS_ASSERT_MOCK
bool "Assert mock for test purposes"
default n
---help---
Invalid parameters make the API return -EFAULT instead of
asserting.
endif # ZBUS

25
system/zbus/Make.defs Normal file
View file

@ -0,0 +1,25 @@
############################################################################
# apps/system/zbus/Make.defs
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership. The
# ASF licenses this file to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance with the
# License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations
# under the License.
#
############################################################################
ifneq ($(CONFIG_ZBUS),)
CONFIGURED_APPS += $(APPDIR)/system/zbus
endif

33
system/zbus/Makefile Normal file
View file

@ -0,0 +1,33 @@
############################################################################
# apps/system/zbus/Makefile
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership. The
# ASF licenses this file to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance with the
# License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations
# under the License.
#
############################################################################
include $(APPDIR)/Make.defs
# ZBus message bus library (Zephyr zbus port)
CSRCS = zbus.c zbus_iterable_sections.c
ifneq ($(CONFIG_ZBUS_RUNTIME_OBSERVERS),)
CSRCS += zbus_runtime_observers.c
endif
include $(APPDIR)/Application.mk

939
system/zbus/zbus.c Normal file
View file

@ -0,0 +1,939 @@
/****************************************************************************
* apps/system/zbus/zbus.c
*
* SPDX-License-Identifier: Apache-2.0
*
* Copyright (c) 2022 Rodrigo Peixoto <rodrigopex@gmail.com>
* Copyright (c) 2026 NuttX port
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may
* not use this file except in compliance with the License. You may obtain
* a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
* License for the specific language governing permissions and limitations
* under the License.
*
****************************************************************************/
/****************************************************************************
* Included Files
****************************************************************************/
#include <nuttx/config.h>
#include <nuttx/clock.h>
#include <nuttx/compiler.h>
#include <nuttx/mqueue.h>
#include <fcntl.h>
#include <inttypes.h>
#include <mqueue.h>
#include <pthread.h>
#include <sched.h>
#include <stdio.h>
#include <stdlib.h>
#include <syslog.h>
#include <system/zbus.h>
#include "zbus_priv.h"
/****************************************************************************
* Private Data
****************************************************************************/
static pthread_once_t g_zbus_once = PTHREAD_ONCE_INIT;
/* Protects observer enabled flags and observation masks */
static pthread_mutex_t g_zbus_obs_lock = PTHREAD_MUTEX_INITIALIZER;
/****************************************************************************
* Private Functions
****************************************************************************/
/****************************************************************************
* Name: zbus_deadline_to_realtime
*
* Description:
* Convert the remaining time of a monotonic deadline into an absolute
* CLOCK_REALTIME timespec as required by mq_timedsend/mq_timedreceive.
*
****************************************************************************/
static void zbus_deadline_to_realtime(FAR const struct zbus_deadline *d,
FAR struct timespec *rt)
{
struct timespec remaining;
struct timespec now;
clock_gettime(CLOCK_REALTIME, rt);
if (d->mode == ZBUS_DEADLINE_ABS)
{
clock_gettime(CLOCK_MONOTONIC, &now);
/* clock_timespec_subtract() returns zero when the deadline has
* already expired, which leaves rt at the current time.
*/
clock_timespec_subtract(&d->abs, &now, &remaining);
clock_timespec_add(rt, &remaining, rt);
}
}
/****************************************************************************
* Name: zbus_mq_send / zbus_mq_recv
*
* Description:
* Message queue send/receive honoring a zbus_deadline. Following the
* Zephyr k_msgq semantics, a no-wait failure returns -ENOMSG and a
* timeout returns -EAGAIN.
*
****************************************************************************/
static int zbus_mq_send(struct file *mq, const char *buf, size_t len,
const struct zbus_deadline *d)
{
struct timespec rt;
int ret;
if (mq->f_inode == NULL)
{
return -ENODEV;
}
if (d->mode == ZBUS_DEADLINE_FOREVER)
{
do
{
ret = file_mq_send(mq, buf, len, 0);
}
while (ret == -EINTR);
}
else
{
zbus_deadline_to_realtime(d, &rt);
do
{
ret = file_mq_timedsend(mq, buf, len, 0, &rt);
}
while (ret == -EINTR);
}
if (ret == -ETIMEDOUT)
{
return (d->mode == ZBUS_DEADLINE_NOWAIT) ? -ENOMSG : -EAGAIN;
}
return ret;
}
static ssize_t zbus_mq_recv(struct file *mq, char *buf, size_t len,
const struct zbus_deadline *d)
{
struct timespec rt;
ssize_t ret;
if (mq->f_inode == NULL)
{
return -ENODEV;
}
if (d->mode == ZBUS_DEADLINE_FOREVER)
{
do
{
ret = file_mq_receive(mq, buf, len, NULL);
}
while (ret == -EINTR);
}
else
{
zbus_deadline_to_realtime(d, &rt);
do
{
ret = file_mq_timedreceive(mq, buf, len, NULL, &rt);
}
while (ret == -EINTR);
}
if (ret == -ETIMEDOUT)
{
return (d->mode == ZBUS_DEADLINE_NOWAIT) ? -ENOMSG : -EAGAIN;
}
return ret;
}
#ifdef CONFIG_ZBUS_ASYNC_LISTENER
/****************************************************************************
* Name: zbus_async_listener_task
*
* Description:
* Dedicated task of an async listener: block on the listener's queue
* and invoke its callback for every message copy, from an aligned
* buffer. argv[1] carries the observer address.
*
****************************************************************************/
static int zbus_async_listener_task(int argc, FAR char *argv[])
{
const struct zbus_observer *obs;
char buf[sizeof(struct zbus_channel *) +
CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE];
uint8_t msg[CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE] aligned_data(8);
const struct zbus_channel *chan;
struct zbus_deadline d;
ssize_t nbytes;
if (argc < 2)
{
return EXIT_FAILURE;
}
obs = (const struct zbus_observer *)(uintptr_t)strtoul(argv[1], NULL, 16);
d.mode = ZBUS_DEADLINE_FOREVER;
for (; ; )
{
nbytes = zbus_mq_recv(&obs->data->mq, buf, sizeof(buf), &d);
if (nbytes < (ssize_t)sizeof(struct zbus_channel *))
{
continue;
}
memcpy(&chan, buf, sizeof(chan));
memcpy(msg, buf + sizeof(chan), nbytes - sizeof(chan));
obs->async_callback(chan, msg);
}
return EXIT_SUCCESS;
}
/****************************************************************************
* Name: zbus_async_listener_start
*
* Description:
* Spawn the task serving an async listener. A task rather than a
* pthread: the lazy init runs in the context of the first API caller,
* and a pthread would die with that caller's task group.
*
****************************************************************************/
static int zbus_async_listener_start(const struct zbus_observer *obs)
{
char arg[2 + sizeof(uintptr_t) * 2 + 1];
FAR char *argv[2];
int pid;
snprintf(arg, sizeof(arg), "%" PRIxPTR, (uintptr_t)obs);
argv[0] = arg;
argv[1] = NULL;
pid = task_create("zbus_async", CONFIG_ZBUS_ASYNC_LISTENER_PRIORITY,
CONFIG_ZBUS_ASYNC_LISTENER_STACKSIZE,
zbus_async_listener_task, argv);
if (pid < 0)
{
return -errno;
}
obs->data->pid = pid;
return 0;
}
#endif /* CONFIG_ZBUS_ASYNC_LISTENER */
/****************************************************************************
* Name: zbus_init_fn
*
* Description:
* One-time initialization: compute the observation index boundaries of
* every channel (relies on the linker sorting the observation section by
* name, which groups entries per channel in priority order) and open the
* notification queues of subscriber-type observers.
*
****************************************************************************/
static void zbus_init_fn(void)
{
FAR struct zbus_channel_observation *observation;
FAR struct zbus_observer *obs;
const struct zbus_channel *curr = NULL;
const struct zbus_channel *prev = NULL;
STRUCT_SECTION_FOREACH(zbus_channel_observation, observation)
{
/* Apply the ROM-preserved initial mask value */
*observation->mask = observation->mask_init;
curr = observation->chan;
if (prev != curr)
{
if (prev == NULL)
{
curr->data->observers_start_idx = 0;
curr->data->observers_end_idx = 0;
}
else
{
curr->data->observers_start_idx =
prev->data->observers_end_idx;
curr->data->observers_end_idx =
prev->data->observers_end_idx;
}
prev = curr;
}
++(curr->data->observers_end_idx);
}
/* Open the notification queues */
STRUCT_SECTION_FOREACH(zbus_observer, obs)
{
struct mq_attr attr;
char name[24];
int ret;
if (obs->type != ZBUS_OBSERVER_SUBSCRIBER_TYPE
#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
&& obs->type != ZBUS_OBSERVER_MSG_SUBSCRIBER_TYPE
#endif
#ifdef CONFIG_ZBUS_ASYNC_LISTENER
&& obs->type != ZBUS_OBSERVER_ASYNC_LISTENER_TYPE
#endif
)
{
continue;
}
memset(&attr, 0, sizeof(attr));
attr.mq_maxmsg = (obs->queue_size > 0) ? obs->queue_size : 1;
#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
if (obs->type != ZBUS_OBSERVER_SUBSCRIBER_TYPE)
{
/* Message subscribers and async listeners carry a copy of the
* message after the channel pointer.
*/
attr.mq_msgsize = sizeof(struct zbus_channel *) +
CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE;
}
else
#endif
{
attr.mq_msgsize = sizeof(struct zbus_channel *);
}
snprintf(name, sizeof(name), "zb%08" PRIxPTR, (uintptr_t)obs);
/* file_mq_open() creates a queue usable from any task, unlike
* mq_open() whose descriptor belongs to the calling task only.
*/
ret = file_mq_open(&obs->data->mq, name, O_RDWR | O_CREAT, 0644,
&attr);
if (ret < 0)
{
syslog(LOG_ERR, "zbus: cannot open queue %s: %d\n", name, ret);
continue;
}
#ifdef CONFIG_ZBUS_ASYNC_LISTENER
if (obs->type == ZBUS_OBSERVER_ASYNC_LISTENER_TYPE)
{
ret = zbus_async_listener_start(obs);
if (ret < 0)
{
syslog(LOG_ERR, "zbus: cannot start async listener %p: %d\n",
obs, ret);
}
}
#endif
}
}
/****************************************************************************
* Name: zbus_notify_observer
*
* Description:
* Deliver one notification. msgbuf carries the pre-built message
* subscriber datagram ({channel pointer, message copy}) or NULL when
* CONFIG_ZBUS_MSG_SUBSCRIBER is disabled.
*
****************************************************************************/
static int zbus_notify_observer(const struct zbus_channel *chan,
const struct zbus_observer *obs,
const struct zbus_deadline *d,
const char *msgbuf)
{
switch (obs->type)
{
case ZBUS_OBSERVER_LISTENER_TYPE:
obs->callback(chan);
return 0;
case ZBUS_OBSERVER_SUBSCRIBER_TYPE:
return zbus_mq_send(&obs->data->mq, (const char *)&chan,
sizeof(chan), d);
#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
case ZBUS_OBSERVER_MSG_SUBSCRIBER_TYPE:
if (chan->message_size > CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE)
{
return -EMSGSIZE;
}
return zbus_mq_send(&obs->data->mq, msgbuf,
sizeof(struct zbus_channel *) +
chan->message_size, d);
#endif
#ifdef CONFIG_ZBUS_ASYNC_LISTENER
case ZBUS_OBSERVER_ASYNC_LISTENER_TYPE:
{
if (chan->message_size > CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE)
{
return -EMSGSIZE;
}
/* The listener's task drains the queue and runs the callback */
return zbus_mq_send(&obs->data->mq, msgbuf,
sizeof(struct zbus_channel *) +
chan->message_size, d);
}
#endif
default:
return -EINVAL;
}
}
/****************************************************************************
* Name: zbus_vded_exec
*
* Description:
* The event dispatcher: notify every enabled/unmasked observer of the
* channel. The channel must be locked by the caller.
*
****************************************************************************/
static int zbus_vded_exec(const struct zbus_channel *chan,
const struct zbus_deadline *d)
{
const char *msgbuf = NULL;
int last_error = 0;
int16_t i;
int err;
#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
char buf[sizeof(struct zbus_channel *) +
CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE];
memcpy(buf, &chan, sizeof(chan));
if (chan->message_size <= CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE)
{
memcpy(buf + sizeof(chan), chan->message, chan->message_size);
}
msgbuf = buf;
#endif
/* The observations of a channel are contiguous and sorted by the linker,
* so notifying them in order follows the priority of the definition.
*/
for (i = chan->data->observers_start_idx;
i < chan->data->observers_end_idx; i++)
{
struct zbus_channel_observation *observation;
const struct zbus_observer *obs;
STRUCT_SECTION_GET(zbus_channel_observation, i, &observation);
obs = observation->obs;
if (!obs->data->enabled || *observation->mask)
{
continue;
}
err = zbus_notify_observer(chan, obs, d, msgbuf);
if (err)
{
last_error = err;
syslog(LOG_ERR, "zbus: could not notify observer %p: %d\n",
obs, err);
}
}
#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
struct zbus_observer_node *obs_nd;
list_for_every_entry(&chan->data->observers, obs_nd,
struct zbus_observer_node, node)
{
if (!obs_nd->obs->data->enabled)
{
continue;
}
err = zbus_notify_observer(chan, obs_nd->obs, d, msgbuf);
if (err)
{
last_error = err;
}
}
#endif
return last_error;
}
/****************************************************************************
* Public Functions
****************************************************************************/
/****************************************************************************
* Name: zbus_port_init_once
****************************************************************************/
void zbus_port_init_once(void)
{
pthread_once(&g_zbus_once, zbus_init_fn);
}
/****************************************************************************
* Name: zbus_deadline_init
****************************************************************************/
void zbus_deadline_init(int32_t timeout_ms, struct zbus_deadline *d)
{
if (timeout_ms < 0)
{
d->mode = ZBUS_DEADLINE_FOREVER;
}
else if (timeout_ms == 0)
{
d->mode = ZBUS_DEADLINE_NOWAIT;
}
else
{
struct timespec delay;
d->mode = ZBUS_DEADLINE_ABS;
clock_gettime(CLOCK_MONOTONIC, &d->abs);
clock_nsec2time(&delay, (int64_t)timeout_ms * NSEC_PER_MSEC);
clock_timespec_add(&d->abs, &delay, &d->abs);
}
}
/****************************************************************************
* Name: zbus_sem_take
****************************************************************************/
int zbus_sem_take(sem_t *sem, const struct zbus_deadline *d)
{
int ret;
switch (d->mode)
{
case ZBUS_DEADLINE_FOREVER:
do
{
ret = sem_wait(sem);
}
while (ret < 0 && errno == EINTR);
return (ret < 0) ? -errno : 0;
case ZBUS_DEADLINE_NOWAIT:
ret = sem_trywait(sem);
if (ret < 0)
{
return (errno == EAGAIN) ? -EBUSY : -errno;
}
return 0;
case ZBUS_DEADLINE_ABS:
default:
do
{
ret = sem_clockwait(sem, CLOCK_MONOTONIC, &d->abs);
}
while (ret < 0 && errno == EINTR);
if (ret < 0)
{
return (errno == ETIMEDOUT) ? -EAGAIN : -errno;
}
return 0;
}
}
/****************************************************************************
* Name: zbus_chan_pub
****************************************************************************/
int zbus_chan_pub(const struct zbus_channel *chan, const void *msg,
int32_t timeout_ms)
{
struct zbus_deadline d;
int err;
_ZBUS_ASSERT(chan != NULL, "chan is required");
_ZBUS_ASSERT(msg != NULL, "msg is required");
zbus_port_init_once();
if (chan->validator != NULL &&
!chan->validator(msg, chan->message_size))
{
return -ENOMSG;
}
zbus_deadline_init(timeout_ms, &d);
err = zbus_sem_take(&chan->data->sem, &d);
if (err)
{
return err;
}
#ifdef CONFIG_ZBUS_CHANNEL_PUBLISH_STATS
zbus_chan_pub_stats_update(chan);
#endif
memcpy(chan->message, msg, chan->message_size);
err = zbus_vded_exec(chan, &d);
sem_post(&chan->data->sem);
return err;
}
/****************************************************************************
* Name: zbus_chan_read
****************************************************************************/
int zbus_chan_read(const struct zbus_channel *chan, void *msg,
int32_t timeout_ms)
{
struct zbus_deadline d;
int err;
_ZBUS_ASSERT(chan != NULL, "chan is required");
_ZBUS_ASSERT(msg != NULL, "msg is required");
zbus_port_init_once();
zbus_deadline_init(timeout_ms, &d);
err = zbus_sem_take(&chan->data->sem, &d);
if (err)
{
return err;
}
memcpy(msg, chan->message, chan->message_size);
sem_post(&chan->data->sem);
return 0;
}
/****************************************************************************
* Name: zbus_chan_notify
****************************************************************************/
int zbus_chan_notify(const struct zbus_channel *chan, int32_t timeout_ms)
{
struct zbus_deadline d;
int err;
_ZBUS_ASSERT(chan != NULL, "chan is required");
zbus_port_init_once();
zbus_deadline_init(timeout_ms, &d);
err = zbus_sem_take(&chan->data->sem, &d);
if (err)
{
return err;
}
err = zbus_vded_exec(chan, &d);
sem_post(&chan->data->sem);
return err;
}
/****************************************************************************
* Name: zbus_chan_claim
****************************************************************************/
int zbus_chan_claim(const struct zbus_channel *chan, int32_t timeout_ms)
{
struct zbus_deadline d;
_ZBUS_ASSERT(chan != NULL, "chan is required");
zbus_port_init_once();
zbus_deadline_init(timeout_ms, &d);
return zbus_sem_take(&chan->data->sem, &d);
}
/****************************************************************************
* Name: zbus_chan_finish
****************************************************************************/
int zbus_chan_finish(const struct zbus_channel *chan)
{
_ZBUS_ASSERT(chan != NULL, "chan is required");
sem_post(&chan->data->sem);
return 0;
}
/****************************************************************************
* Name: zbus_sub_wait
****************************************************************************/
int zbus_sub_wait(const struct zbus_observer *sub,
const struct zbus_channel **chan, int32_t timeout_ms)
{
const struct zbus_channel *received;
struct zbus_deadline d;
ssize_t nbytes;
_ZBUS_ASSERT(sub != NULL, "sub is required");
_ZBUS_ASSERT(sub->type == ZBUS_OBSERVER_SUBSCRIBER_TYPE,
"sub must be a SUBSCRIBER");
_ZBUS_ASSERT(chan != NULL, "chan is required");
zbus_port_init_once();
zbus_deadline_init(timeout_ms, &d);
nbytes = zbus_mq_recv(&sub->data->mq, (char *)&received,
sizeof(received), &d);
if (nbytes < 0)
{
return nbytes;
}
*chan = received;
return 0;
}
#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
/****************************************************************************
* Name: zbus_sub_wait_msg
****************************************************************************/
int zbus_sub_wait_msg(const struct zbus_observer *sub,
const struct zbus_channel **chan, void *msg,
int32_t timeout_ms)
{
char buf[sizeof(struct zbus_channel *) +
CONFIG_ZBUS_MSG_SUBSCRIBER_MAX_MSG_SIZE];
struct zbus_deadline d;
ssize_t nbytes;
_ZBUS_ASSERT(sub != NULL, "sub is required");
_ZBUS_ASSERT(sub->type == ZBUS_OBSERVER_MSG_SUBSCRIBER_TYPE,
"sub must be a MSG_SUBSCRIBER");
_ZBUS_ASSERT(chan != NULL, "chan is required");
_ZBUS_ASSERT(msg != NULL, "msg is required");
zbus_port_init_once();
zbus_deadline_init(timeout_ms, &d);
nbytes = zbus_mq_recv(&sub->data->mq, buf, sizeof(buf), &d);
if (nbytes < 0)
{
return nbytes;
}
if (nbytes < (ssize_t)sizeof(struct zbus_channel *))
{
return -EILSEQ;
}
memcpy(chan, buf, sizeof(struct zbus_channel *));
memcpy(msg, buf + sizeof(struct zbus_channel *),
nbytes - sizeof(struct zbus_channel *));
return 0;
}
#endif /* CONFIG_ZBUS_MSG_SUBSCRIBER */
/****************************************************************************
* Name: zbus_obs_set_enable
****************************************************************************/
int zbus_obs_set_enable(const struct zbus_observer *obs, bool enabled)
{
_ZBUS_ASSERT(obs != NULL, "obs is required");
pthread_mutex_lock(&g_zbus_obs_lock);
obs->data->enabled = enabled;
pthread_mutex_unlock(&g_zbus_obs_lock);
return 0;
}
/****************************************************************************
* Name: zbus_obs_set_chan_notification_mask
****************************************************************************/
int zbus_obs_set_chan_notification_mask(const struct zbus_observer *obs,
const struct zbus_channel *chan,
bool masked)
{
int16_t i;
int err = -ESRCH;
_ZBUS_ASSERT(obs != NULL, "obs is required");
_ZBUS_ASSERT(chan != NULL, "chan is required");
zbus_port_init_once();
pthread_mutex_lock(&g_zbus_obs_lock);
for (i = chan->data->observers_start_idx;
i < chan->data->observers_end_idx; i++)
{
struct zbus_channel_observation *observation;
STRUCT_SECTION_GET(zbus_channel_observation, i, &observation);
if (observation->obs == obs)
{
*observation->mask = masked;
err = 0;
break;
}
}
pthread_mutex_unlock(&g_zbus_obs_lock);
return err;
}
/****************************************************************************
* Name: zbus_obs_is_chan_notification_masked
****************************************************************************/
int zbus_obs_is_chan_notification_masked(const struct zbus_observer *obs,
const struct zbus_channel *chan,
bool *masked)
{
int16_t i;
int err = -ESRCH;
_ZBUS_ASSERT(obs != NULL, "obs is required");
_ZBUS_ASSERT(chan != NULL, "chan is required");
_ZBUS_ASSERT(masked != NULL, "masked is required");
zbus_port_init_once();
pthread_mutex_lock(&g_zbus_obs_lock);
for (i = chan->data->observers_start_idx;
i < chan->data->observers_end_idx; i++)
{
struct zbus_channel_observation *observation;
STRUCT_SECTION_GET(zbus_channel_observation, i, &observation);
if (observation->obs == obs)
{
*masked = *observation->mask;
err = 0;
break;
}
}
pthread_mutex_unlock(&g_zbus_obs_lock);
return err;
}
#ifdef CONFIG_ZBUS_CHANNEL_ID
/****************************************************************************
* Name: zbus_chan_from_id
****************************************************************************/
const struct zbus_channel *zbus_chan_from_id(uint32_t channel_id)
{
FAR struct zbus_channel *chan;
if (channel_id == ZBUS_CHAN_ID_INVALID)
{
return NULL;
}
STRUCT_SECTION_FOREACH(zbus_channel, chan)
{
if (chan->id == channel_id)
{
return chan;
}
}
return NULL;
}
#endif /* CONFIG_ZBUS_CHANNEL_ID */
#ifdef CONFIG_ZBUS_CHANNEL_NAME
/****************************************************************************
* Name: zbus_chan_from_name
****************************************************************************/
const struct zbus_channel *zbus_chan_from_name(const char *name)
{
FAR struct zbus_channel *chan;
if (name == NULL)
{
return NULL;
}
STRUCT_SECTION_FOREACH(zbus_channel, chan)
{
if (strcmp(chan->name, name) == 0)
{
return chan;
}
}
return NULL;
}
#endif /* CONFIG_ZBUS_CHANNEL_NAME */

View file

@ -0,0 +1,101 @@
/****************************************************************************
* apps/system/zbus/zbus_iterable_sections.c
*
* SPDX-License-Identifier: Apache-2.0
*
* Copyright (c) 2022 Rodrigo Peixoto <rodrigopex@gmail.com>
* Copyright (c) 2026 NuttX port
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may
* not use this file except in compliance with the License. You may obtain
* a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
* License for the specific language governing permissions and limitations
* under the License.
*
****************************************************************************/
/****************************************************************************
* Included Files
****************************************************************************/
#include <nuttx/config.h>
#include <system/zbus.h>
#include "zbus_priv.h"
/****************************************************************************
* Public Functions
****************************************************************************/
bool zbus_iterate_over_channels(
bool (*iterator_func)(const struct zbus_channel *chan))
{
FAR struct zbus_channel *chan;
STRUCT_SECTION_FOREACH(zbus_channel, chan)
{
if (!(*iterator_func)(chan))
{
return false;
}
}
return true;
}
bool zbus_iterate_over_channels_with_user_data(
bool (*iterator_func)(const struct zbus_channel *chan, void *user_data),
void *user_data)
{
FAR struct zbus_channel *chan;
STRUCT_SECTION_FOREACH(zbus_channel, chan)
{
if (!(*iterator_func)(chan, user_data))
{
return false;
}
}
return true;
}
bool zbus_iterate_over_observers(
bool (*iterator_func)(const struct zbus_observer *obs))
{
FAR struct zbus_observer *obs;
STRUCT_SECTION_FOREACH(zbus_observer, obs)
{
if (!(*iterator_func)(obs))
{
return false;
}
}
return true;
}
bool zbus_iterate_over_observers_with_user_data(
bool (*iterator_func)(const struct zbus_observer *obs, void *user_data),
void *user_data)
{
FAR struct zbus_observer *obs;
STRUCT_SECTION_FOREACH(zbus_observer, obs)
{
if (!(*iterator_func)(obs, user_data))
{
return false;
}
}
return true;
}

85
system/zbus/zbus_priv.h Normal file
View file

@ -0,0 +1,85 @@
/****************************************************************************
* apps/system/zbus/zbus_priv.h
*
* SPDX-License-Identifier: Apache-2.0
*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership. The
* ASF licenses this file to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance with the
* License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
* License for the specific language governing permissions and limitations
* under the License.
*
****************************************************************************/
#ifndef __APPS_SYSTEM_ZBUS_ZBUS_PRIV_H
#define __APPS_SYSTEM_ZBUS_ZBUS_PRIV_H
/****************************************************************************
* Included Files
****************************************************************************/
#include <nuttx/config.h>
#include <semaphore.h>
#include <stdint.h>
#include <time.h>
#include <system/zbus.h>
/****************************************************************************
* Public Data
****************************************************************************/
/* Boundary symbols of the zbus iterable sections */
STRUCT_SECTION_DECLARE(zbus_channel);
STRUCT_SECTION_DECLARE(zbus_observer);
STRUCT_SECTION_DECLARE(zbus_channel_observation);
/****************************************************************************
* Public Types
****************************************************************************/
/* Deadline computed once per API call and honored by every internal wait */
enum zbus_deadline_mode_e
{
ZBUS_DEADLINE_FOREVER = 0,
ZBUS_DEADLINE_NOWAIT,
ZBUS_DEADLINE_ABS
};
struct zbus_deadline
{
enum zbus_deadline_mode_e mode;
struct timespec abs; /* CLOCK_MONOTONIC absolute deadline */
};
/****************************************************************************
* Public Function Prototypes
****************************************************************************/
/* One-time lazy initialization (observation indexes, observer queues) */
void zbus_port_init_once(void);
/* Deadline helpers */
void zbus_deadline_init(int32_t timeout_ms, struct zbus_deadline *d);
/* Take a semaphore honoring the deadline. Returns 0, -EBUSY (no-wait) or
* -EAGAIN (timed out).
*/
int zbus_sem_take(sem_t *sem, const struct zbus_deadline *d);
#endif /* __APPS_SYSTEM_ZBUS_ZBUS_PRIV_H */

View file

@ -0,0 +1,147 @@
/****************************************************************************
* apps/system/zbus/zbus_runtime_observers.c
*
* SPDX-License-Identifier: Apache-2.0
*
* Copyright (c) 2022 Rodrigo Peixoto <rodrigopex@gmail.com>
* Copyright (c) 2026 NuttX port
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may
* not use this file except in compliance with the License. You may obtain
* a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
* License for the specific language governing permissions and limitations
* under the License.
*
****************************************************************************/
/****************************************************************************
* Included Files
****************************************************************************/
#include <nuttx/config.h>
#include <stdlib.h>
#include <system/zbus.h>
#include "zbus_priv.h"
/****************************************************************************
* Public Functions
****************************************************************************/
/****************************************************************************
* Name: zbus_chan_add_obs
****************************************************************************/
int zbus_chan_add_obs(const struct zbus_channel *chan,
const struct zbus_observer *obs, int32_t timeout_ms)
{
struct zbus_observer_node *obs_nd;
struct zbus_deadline d;
int err;
_ZBUS_ASSERT(chan != NULL, "chan is required");
_ZBUS_ASSERT(obs != NULL, "obs is required");
zbus_port_init_once();
zbus_deadline_init(timeout_ms, &d);
err = zbus_sem_take(&chan->data->sem, &d);
if (err)
{
return err;
}
/* Reject observers already statically attached to the channel */
for (int16_t i = chan->data->observers_start_idx,
limit = chan->data->observers_end_idx; i < limit; i++)
{
struct zbus_channel_observation *observation;
STRUCT_SECTION_GET(zbus_channel_observation, i, &observation);
if (observation->obs == obs)
{
sem_post(&chan->data->sem);
return -EEXIST;
}
}
/* Reject observers already dynamically attached to the channel */
list_for_every_entry(&chan->data->observers, obs_nd,
struct zbus_observer_node, node)
{
if (obs_nd->obs == obs)
{
sem_post(&chan->data->sem);
return -EALREADY;
}
}
obs_nd = malloc(sizeof(*obs_nd));
if (obs_nd == NULL)
{
sem_post(&chan->data->sem);
return -ENOMEM;
}
obs_nd->obs = obs;
list_add_tail(&chan->data->observers, &obs_nd->node);
sem_post(&chan->data->sem);
return 0;
}
/****************************************************************************
* Name: zbus_chan_rm_obs
****************************************************************************/
int zbus_chan_rm_obs(const struct zbus_channel *chan,
const struct zbus_observer *obs, int32_t timeout_ms)
{
struct zbus_observer_node *obs_nd;
struct zbus_observer_node *tmp;
struct zbus_deadline d;
int err;
_ZBUS_ASSERT(chan != NULL, "chan is required");
_ZBUS_ASSERT(obs != NULL, "obs is required");
zbus_port_init_once();
zbus_deadline_init(timeout_ms, &d);
err = zbus_sem_take(&chan->data->sem, &d);
if (err)
{
return err;
}
list_for_every_entry_safe(&chan->data->observers, obs_nd, tmp,
struct zbus_observer_node, node)
{
if (obs_nd->obs == obs)
{
list_delete(&obs_nd->node);
free(obs_nd);
sem_post(&chan->data->sem);
return 0;
}
}
sem_post(&chan->data->sem);
return -ENODATA;
}

View file

@ -0,0 +1,33 @@
# ##############################################################################
# apps/testing/zbus/CMakeLists.txt
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more contributor
# license agreements. See the NOTICE file distributed with this work for
# additional information regarding copyright ownership. The ASF licenses this
# file to you under the Apache License, Version 2.0 (the "License"); you may not
# use this file except in compliance with the License. You may obtain a copy of
# the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations under
# the License.
#
# ##############################################################################
if(CONFIG_TESTING_ZBUS)
nuttx_add_application(
NAME
cmocka_zbus_test
SRCS
zbustest.c
STACKSIZE
${CONFIG_TESTING_ZBUS_STACKSIZE}
PRIORITY
${CONFIG_TESTING_ZBUS_PRIORITY})
endif()

26
testing/zbus/Kconfig Normal file
View file

@ -0,0 +1,26 @@
#
# For a description of the syntax of this configuration file,
# see the file kconfig-language.txt in the NuttX tools repository.
#
config TESTING_ZBUS
tristate "cmocka zbus test"
default n
depends on ZBUS && TESTING_CMOCKA
---help---
Enable the cmocka zbus message bus test suite. Covers channel
publish/read, listeners, subscribers, message subscribers,
validators, notification masks, observer enable/disable,
runtime observers, claim/finish, timeouts and iteration.
if TESTING_ZBUS
config TESTING_ZBUS_PRIORITY
int "zbus test task priority"
default 100
config TESTING_ZBUS_STACKSIZE
int "zbus test stack size"
default 8192
endif # TESTING_ZBUS

25
testing/zbus/Make.defs Normal file
View file

@ -0,0 +1,25 @@
############################################################################
# apps/testing/zbus/Make.defs
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership. The
# ASF licenses this file to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance with the
# License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations
# under the License.
#
############################################################################
ifneq ($(CONFIG_TESTING_ZBUS),)
CONFIGURED_APPS += $(APPDIR)/testing/zbus
endif

34
testing/zbus/Makefile Normal file
View file

@ -0,0 +1,34 @@
############################################################################
# apps/testing/zbus/Makefile
#
# SPDX-License-Identifier: Apache-2.0
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership. The
# ASF licenses this file to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance with the
# License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
# License for the specific language governing permissions and limitations
# under the License.
#
############################################################################
include $(APPDIR)/Make.defs
# cmocka zbus test
PROGNAME = cmocka_zbus_test
PRIORITY = $(CONFIG_TESTING_ZBUS_PRIORITY)
STACKSIZE = $(CONFIG_TESTING_ZBUS_STACKSIZE)
MODULE = $(CONFIG_TESTING_ZBUS)
MAINSRC = zbustest.c
include $(APPDIR)/Application.mk

1028
testing/zbus/zbustest.c Normal file

File diff suppressed because it is too large Load diff