This is an automated email from the ASF dual-hosted git repository.

acassis pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/nuttx-apps.git


The following commit(s) were added to refs/heads/master by this push:
     new ee20ddd0b system/zbus: Port the Zephyr zbus message bus to NuttX
ee20ddd0b is described below

commit ee20ddd0ba47dedd2d66949d1d869b22065c21ec
Author: Jorge Guzman <[email protected]>
AuthorDate: Fri Jul 10 09:23:58 2026 -0300

    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 <[email protected]>
---
 examples/zbus/CMakeLists.txt         |   33 ++
 examples/zbus/Kconfig                |   27 +
 examples/zbus/Make.defs              |   25 +
 examples/zbus/Makefile               |   34 ++
 examples/zbus/zbus_main.c            |  160 ++++++
 include/system/zbus.h                |  719 ++++++++++++++++++++++++
 system/zbus/CMakeLists.txt           |   31 +
 system/zbus/Kconfig                  |  110 ++++
 system/zbus/Make.defs                |   25 +
 system/zbus/Makefile                 |   33 ++
 system/zbus/zbus.c                   |  939 +++++++++++++++++++++++++++++++
 system/zbus/zbus_iterable_sections.c |  101 ++++
 system/zbus/zbus_priv.h              |   85 +++
 system/zbus/zbus_runtime_observers.c |  147 +++++
 testing/zbus/CMakeLists.txt          |   33 ++
 testing/zbus/Kconfig                 |   26 +
 testing/zbus/Make.defs               |   25 +
 testing/zbus/Makefile                |   34 ++
 testing/zbus/zbustest.c              | 1028 ++++++++++++++++++++++++++++++++++
 19 files changed, 3615 insertions(+)

diff --git a/examples/zbus/CMakeLists.txt b/examples/zbus/CMakeLists.txt
new file mode 100644
index 000000000..8c476c4b6
--- /dev/null
+++ b/examples/zbus/CMakeLists.txt
@@ -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()
diff --git a/examples/zbus/Kconfig b/examples/zbus/Kconfig
new file mode 100644
index 000000000..fe0fd77f6
--- /dev/null
+++ b/examples/zbus/Kconfig
@@ -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
diff --git a/examples/zbus/Make.defs b/examples/zbus/Make.defs
new file mode 100644
index 000000000..d21f2d347
--- /dev/null
+++ b/examples/zbus/Make.defs
@@ -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
diff --git a/examples/zbus/Makefile b/examples/zbus/Makefile
new file mode 100644
index 000000000..b194bdfb0
--- /dev/null
+++ b/examples/zbus/Makefile
@@ -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
diff --git a/examples/zbus/zbus_main.c b/examples/zbus/zbus_main.c
new file mode 100644
index 000000000..a37ee4414
--- /dev/null
+++ b/examples/zbus/zbus_main.c
@@ -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;
+}
diff --git a/include/system/zbus.h b/include/system/zbus.h
new file mode 100644
index 000000000..a729e2f41
--- /dev/null
+++ b/include/system/zbus.h
@@ -0,0 +1,719 @@
+/****************************************************************************
+ * apps/include/system/zbus.h
+ *
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * Copyright (c) 2022 Rodrigo Peixoto <[email protected]>
+ * 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 */
diff --git a/system/zbus/CMakeLists.txt b/system/zbus/CMakeLists.txt
new file mode 100644
index 000000000..900977750
--- /dev/null
+++ b/system/zbus/CMakeLists.txt
@@ -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()
diff --git a/system/zbus/Kconfig b/system/zbus/Kconfig
new file mode 100644
index 000000000..2fe59468b
--- /dev/null
+++ b/system/zbus/Kconfig
@@ -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
diff --git a/system/zbus/Make.defs b/system/zbus/Make.defs
new file mode 100644
index 000000000..4b4ac4d3a
--- /dev/null
+++ b/system/zbus/Make.defs
@@ -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
diff --git a/system/zbus/Makefile b/system/zbus/Makefile
new file mode 100644
index 000000000..c136a27cf
--- /dev/null
+++ b/system/zbus/Makefile
@@ -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
diff --git a/system/zbus/zbus.c b/system/zbus/zbus.c
new file mode 100644
index 000000000..124d4ee97
--- /dev/null
+++ b/system/zbus/zbus.c
@@ -0,0 +1,939 @@
+/****************************************************************************
+ * apps/system/zbus/zbus.c
+ *
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * Copyright (c) 2022 Rodrigo Peixoto <[email protected]>
+ * 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 */
diff --git a/system/zbus/zbus_iterable_sections.c 
b/system/zbus/zbus_iterable_sections.c
new file mode 100644
index 000000000..1fda7b9f3
--- /dev/null
+++ b/system/zbus/zbus_iterable_sections.c
@@ -0,0 +1,101 @@
+/****************************************************************************
+ * apps/system/zbus/zbus_iterable_sections.c
+ *
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * Copyright (c) 2022 Rodrigo Peixoto <[email protected]>
+ * 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;
+}
diff --git a/system/zbus/zbus_priv.h b/system/zbus/zbus_priv.h
new file mode 100644
index 000000000..8ffa89872
--- /dev/null
+++ b/system/zbus/zbus_priv.h
@@ -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 */
diff --git a/system/zbus/zbus_runtime_observers.c 
b/system/zbus/zbus_runtime_observers.c
new file mode 100644
index 000000000..9877c65b5
--- /dev/null
+++ b/system/zbus/zbus_runtime_observers.c
@@ -0,0 +1,147 @@
+/****************************************************************************
+ * apps/system/zbus/zbus_runtime_observers.c
+ *
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * Copyright (c) 2022 Rodrigo Peixoto <[email protected]>
+ * 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;
+}
diff --git a/testing/zbus/CMakeLists.txt b/testing/zbus/CMakeLists.txt
new file mode 100644
index 000000000..bb14401e1
--- /dev/null
+++ b/testing/zbus/CMakeLists.txt
@@ -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()
diff --git a/testing/zbus/Kconfig b/testing/zbus/Kconfig
new file mode 100644
index 000000000..fb7b752c7
--- /dev/null
+++ b/testing/zbus/Kconfig
@@ -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
diff --git a/testing/zbus/Make.defs b/testing/zbus/Make.defs
new file mode 100644
index 000000000..6c6cdda71
--- /dev/null
+++ b/testing/zbus/Make.defs
@@ -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
diff --git a/testing/zbus/Makefile b/testing/zbus/Makefile
new file mode 100644
index 000000000..b9a0039b1
--- /dev/null
+++ b/testing/zbus/Makefile
@@ -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
diff --git a/testing/zbus/zbustest.c b/testing/zbus/zbustest.c
new file mode 100644
index 000000000..5f63a1a34
--- /dev/null
+++ b/testing/zbus/zbustest.c
@@ -0,0 +1,1028 @@
+/****************************************************************************
+ * apps/testing/zbus/zbustest.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 <assert.h>
+#include <math.h>
+#include <pthread.h>
+#include <setjmp.h>
+#include <signal.h>
+#include <stdarg.h>
+#include <stddef.h>
+#include <stdint.h>
+#include <time.h>
+#include <unistd.h>
+#include <cmocka.h>
+
+#include <system/zbus.h>
+
+/****************************************************************************
+ * Private Types
+ ****************************************************************************/
+
+struct zbt_msg_s
+{
+  uint32_t seq;
+  uint32_t value;
+};
+
+/* Sensor-like payload: exercises float/double members through the whole
+ * pipeline (channel storage, listener const view, read-back and the
+ * message subscriber mq copy), which is the typical zbus use case.
+ */
+
+struct zbt_imu_msg_s
+{
+  float accel[3];
+  double magnitude;
+  uint32_t seq;
+};
+
+/****************************************************************************
+ * Private Data
+ ****************************************************************************/
+
+static int g_listener_a_count;
+static uint32_t g_listener_a_last;
+static int g_listener_b_count;
+static uint32_t g_listener_b_last;
+static int g_listener_rt_count;
+
+#ifdef CONFIG_ZBUS_ASYNC_LISTENER
+static volatile int g_async_count;
+static volatile uint32_t g_async_last;
+#endif
+
+static uint32_t g_user_word = 0xcafe;
+
+/****************************************************************************
+ * Private Functions
+ ****************************************************************************/
+
+static void listener_a_cb(const struct zbus_channel *chan)
+{
+  const struct zbt_msg_s *msg = zbus_chan_const_msg(chan);
+
+  g_listener_a_count++;
+  g_listener_a_last = msg->value;
+}
+
+static void listener_b_cb(const struct zbus_channel *chan)
+{
+  const struct zbt_msg_s *msg = zbus_chan_const_msg(chan);
+
+  g_listener_b_count++;
+  g_listener_b_last = msg->value;
+}
+
+static void listener_rt_cb(const struct zbus_channel *chan)
+{
+  (void)chan;
+  g_listener_rt_count++;
+}
+
+#ifdef CONFIG_ZBUS_ASYNC_LISTENER
+static void async_listener_cb(const struct zbus_channel *chan,
+                              const void *msg)
+{
+  const struct zbt_msg_s *m = msg;
+
+  (void)chan;
+  g_async_last = m->value;
+  g_async_count++;
+}
+#endif
+
+/* Notification order: each callback records when it ran, so the test can
+ * check that the observers are notified in the order they were listed.
+ */
+
+static int g_order_seq;
+static int g_order_first;
+static int g_order_second;
+static int g_order_added;
+
+static void order_first_cb(const struct zbus_channel *chan)
+{
+  (void)chan;
+  g_order_first = ++g_order_seq;
+}
+
+static void order_second_cb(const struct zbus_channel *chan)
+{
+  (void)chan;
+  g_order_second = ++g_order_seq;
+}
+
+static void order_added_cb(const struct zbus_channel *chan)
+{
+  (void)chan;
+  g_order_added = ++g_order_seq;
+}
+
+/* Validator: rejects messages with value == 0xdead */
+
+static bool chan_b_validator(const void *msg, size_t msg_size)
+{
+  const struct zbt_msg_s *m = msg;
+
+  (void)msg_size;
+  return m->value != 0xdead;
+}
+
+static int g_imu_listener_count;
+static struct zbt_imu_msg_s g_imu_listener_last;
+
+/* Timer-driven sampler (see test_timer_driven_publisher) */
+
+#define ZBT_SAMPLER_SIGNAL   SIGUSR1
+#define ZBT_SAMPLER_SAMPLES  5
+#define ZBT_SAMPLER_PERIOD   10000000L    /* 10 ms */
+
+static volatile int g_sampler_published;
+static volatile int g_sampler_errors;
+
+static void imu_listener_cb(const struct zbus_channel *chan)
+{
+  const struct zbt_imu_msg_s *msg = zbus_chan_const_msg(chan);
+
+  g_imu_listener_last = *msg;
+  g_imu_listener_count++;
+}
+
+/* Validator doing float math: rejects non-finite samples */
+
+static bool imu_validator(const void *msg, size_t msg_size)
+{
+  const struct zbt_imu_msg_s *m = msg;
+
+  (void)msg_size;
+  return isfinite(m->accel[0]) && isfinite(m->accel[1]) &&
+         isfinite(m->accel[2]) && isfinite(m->magnitude);
+}
+
+/* Interrupt-driven sampling, the NuttX way for a userspace library: the
+ * kernel timer interrupt delivers a signal, a thread waits for it and
+ * publishes from thread context.  Sample values are exactly
+ * representable so the consumers can compare them bit-exact.
+ */
+
+ZBUS_CHAN_DECLARE(zbt_chan_imu);
+
+static FAR void *sampler_thread(FAR void *arg)
+{
+  struct zbt_imu_msg_s msg;
+  struct itimerspec its;
+  struct sigevent sev;
+  sigset_t set;
+  timer_t timer;
+  int n = 0;
+
+  (void)arg;
+
+  sigemptyset(&set);
+  sigaddset(&set, ZBT_SAMPLER_SIGNAL);
+  pthread_sigmask(SIG_BLOCK, &set, NULL);
+
+  memset(&sev, 0, sizeof(sev));
+  sev.sigev_notify = SIGEV_SIGNAL;
+  sev.sigev_signo = ZBT_SAMPLER_SIGNAL;
+  if (timer_create(CLOCK_MONOTONIC, &sev, &timer) != 0)
+    {
+      g_sampler_errors++;
+      return NULL;
+    }
+
+  its.it_value.tv_sec = 0;
+  its.it_value.tv_nsec = ZBT_SAMPLER_PERIOD;
+  its.it_interval = its.it_value;
+  timer_settime(timer, 0, &its, NULL);
+
+  while (n < ZBT_SAMPLER_SAMPLES)
+    {
+      if (sigwaitinfo(&set, NULL) < 0)
+        {
+          if (errno == EINTR)
+            {
+              continue;
+            }
+
+          g_sampler_errors++;
+          break;
+        }
+
+      /* Timer tick: "read the sensor" and publish */
+
+      msg.accel[0] = 0.5f * n;
+      msg.accel[1] = -9.5f;
+      msg.accel[2] = 0.25f;
+      msg.magnitude = 9.5 + n;
+      msg.seq = ++n;
+
+      if (zbus_chan_pub(&zbt_chan_imu, &msg, 100) == 0)
+        {
+          g_sampler_published++;
+        }
+      else
+        {
+          g_sampler_errors++;
+        }
+    }
+
+  timer_delete(timer);
+  return NULL;
+}
+
+/****************************************************************************
+ * Channel and observer definitions
+ ****************************************************************************/
+
+ZBUS_LISTENER_DEFINE(zbt_listener_a, listener_a_cb);
+ZBUS_SUBSCRIBER_DEFINE(zbt_sub_a, 4);
+ZBUS_LISTENER_DEFINE(zbt_listener_b, listener_b_cb);
+ZBUS_LISTENER_DEFINE(zbt_listener_rt, listener_rt_cb);
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+ZBUS_MSG_SUBSCRIBER_DEFINE(zbt_msgsub_b);
+#endif
+
+ZBUS_CHAN_DEFINE(zbt_chan_a,
+                 struct zbt_msg_s,
+                 NULL,
+                 NULL,
+                 ZBUS_OBSERVERS(zbt_listener_a, zbt_sub_a),
+                 ZBUS_MSG_INIT(.seq = 0, .value = 0));
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+ZBUS_CHAN_DEFINE(zbt_chan_b,
+                 struct zbt_msg_s,
+                 chan_b_validator,
+                 &g_user_word,
+                 ZBUS_OBSERVERS(zbt_listener_b, zbt_msgsub_b),
+                 ZBUS_MSG_INIT(.seq = 0, .value = 0));
+#else
+ZBUS_CHAN_DEFINE(zbt_chan_b,
+                 struct zbt_msg_s,
+                 chan_b_validator,
+                 &g_user_word,
+                 ZBUS_OBSERVERS(zbt_listener_b),
+                 ZBUS_MSG_INIT(.seq = 0, .value = 0));
+#endif
+
+#ifdef CONFIG_ZBUS_ASYNC_LISTENER
+/* Attached to zbt_chan_b through ZBUS_CHAN_ADD_OBS (also exercises the
+ * out-of-line observation macro).
+ */
+
+ZBUS_ASYNC_LISTENER_DEFINE(zbt_async_l, async_listener_cb);
+ZBUS_CHAN_ADD_OBS(zbt_chan_b, zbt_async_l, 01);
+#endif
+
+ZBUS_LISTENER_DEFINE(zbt_imu_listener, imu_listener_cb);
+
+ZBUS_LISTENER_DEFINE(zbt_ord_first, order_first_cb);
+ZBUS_LISTENER_DEFINE(zbt_ord_second, order_second_cb);
+ZBUS_LISTENER_DEFINE(zbt_ord_added, order_added_cb);
+
+ZBUS_CHAN_DEFINE(zbt_chan_ord,
+                 struct zbt_msg_s,
+                 NULL,
+                 NULL,
+                 ZBUS_OBSERVERS(zbt_ord_first, zbt_ord_second),
+                 ZBUS_MSG_INIT(.seq = 0, .value = 0));
+
+/* Bound out of line:  must be notified after the two listed above */
+
+ZBUS_CHAN_ADD_OBS(zbt_chan_ord, zbt_ord_added, 01);
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+ZBUS_MSG_SUBSCRIBER_DEFINE(zbt_imu_msgsub);
+
+ZBUS_CHAN_DEFINE(zbt_chan_imu,
+                 struct zbt_imu_msg_s,
+                 imu_validator,
+                 NULL,
+                 ZBUS_OBSERVERS(zbt_imu_listener, zbt_imu_msgsub),
+                 ZBUS_MSG_INIT(.seq = 0));
+#else
+ZBUS_CHAN_DEFINE(zbt_chan_imu,
+                 struct zbt_imu_msg_s,
+                 imu_validator,
+                 NULL,
+                 ZBUS_OBSERVERS(zbt_imu_listener),
+                 ZBUS_MSG_INIT(.seq = 0));
+#endif
+
+#ifdef CONFIG_ZBUS_CHANNEL_ID
+ZBUS_CHAN_DEFINE_WITH_ID(zbt_chan_c,
+                         42,
+                         struct zbt_msg_s,
+                         NULL,
+                         NULL,
+                         ZBUS_OBSERVERS_EMPTY,
+                         ZBUS_MSG_INIT(.seq = 0, .value = 0));
+#else
+ZBUS_CHAN_DEFINE(zbt_chan_c,
+                 struct zbt_msg_s,
+                 NULL,
+                 NULL,
+                 ZBUS_OBSERVERS_EMPTY,
+                 ZBUS_MSG_INIT(.seq = 0, .value = 0));
+#endif
+
+/****************************************************************************
+ * Test helpers
+ ****************************************************************************/
+
+static void drain_subscriber(const struct zbus_observer *sub)
+{
+  const struct zbus_channel *chan;
+
+  while (zbus_sub_wait(sub, &chan, ZBUS_NO_WAIT) == 0)
+    {
+    }
+}
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+static void drain_msg_subscriber(const struct zbus_observer *sub)
+{
+  const struct zbus_channel *chan;
+  struct zbt_msg_s msg;
+
+  while (zbus_sub_wait_msg(sub, &chan, &msg, ZBUS_NO_WAIT) == 0)
+    {
+    }
+}
+#endif
+
+static void reset_all(void)
+{
+  g_listener_a_count = 0;
+  g_listener_b_count = 0;
+  g_listener_rt_count = 0;
+
+  zbus_obs_set_enable(&zbt_listener_a, true);
+  zbus_obs_set_enable(&zbt_listener_b, true);
+  zbus_obs_set_chan_notification_mask(&zbt_listener_a, &zbt_chan_a, false);
+
+  drain_subscriber(&zbt_sub_a);
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+  drain_msg_subscriber(&zbt_msgsub_b);
+#endif
+
+#ifdef CONFIG_ZBUS_ASYNC_LISTENER
+  /* Let the async listener task drain deliveries from previous tests
+   * before resetting its counters.
+   */
+
+  usleep(20 * 1000);
+  g_async_count = 0;
+#endif
+}
+
+/****************************************************************************
+ * Test cases
+ ****************************************************************************/
+
+/* Basic publish: listener receives synchronously, subscriber gets the
+ * notification through its queue, read returns the published message.
+ */
+
+static void test_pub_read_listener_subscriber(FAR void **state)
+{
+  const struct zbus_channel *chan;
+  struct zbt_msg_s msg;
+  int ret;
+
+  (void)state;
+  reset_all();
+
+  msg.seq = 1;
+  msg.value = 100;
+  ret = zbus_chan_pub(&zbt_chan_a, &msg, 1000);
+  assert_int_equal(ret, 0);
+
+  assert_int_equal(g_listener_a_count, 1);
+  assert_int_equal(g_listener_a_last, 100);
+
+  ret = zbus_sub_wait(&zbt_sub_a, &chan, 1000);
+  assert_int_equal(ret, 0);
+  assert_ptr_equal(chan, &zbt_chan_a);
+
+  memset(&msg, 0, sizeof(msg));
+  ret = zbus_chan_read(&zbt_chan_a, &msg, 500);
+  assert_int_equal(ret, 0);
+  assert_int_equal(msg.value, 100);
+}
+
+/* Publishing to one channel must not notify observers of another channel
+ * (validates the observation index grouping computed at init).
+ */
+
+static void test_multi_channel_isolation(FAR void **state)
+{
+  struct zbt_msg_s msg;
+
+  (void)state;
+  reset_all();
+
+  msg.seq = 1;
+  msg.value = 111;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, 1000), 0);
+  assert_int_equal(g_listener_a_count, 1);
+  assert_int_equal(g_listener_b_count, 0);
+
+  msg.value = 222;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_b, &msg, 1000), 0);
+  assert_int_equal(g_listener_b_count, 1);
+  assert_int_equal(g_listener_b_last, 222);
+  assert_int_equal(g_listener_a_count, 1);
+
+  /* Channel with no observers: publish must succeed and reach nobody */
+
+  msg.value = 333;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_c, &msg, 1000), 0);
+  assert_int_equal(g_listener_a_count, 1);
+  assert_int_equal(g_listener_b_count, 1);
+}
+
+/* Validator: invalid messages are rejected with -ENOMSG and nobody is
+ * notified.
+ */
+
+static void test_validator(FAR void **state)
+{
+  struct zbt_msg_s msg;
+
+  (void)state;
+  reset_all();
+
+  msg.seq = 1;
+  msg.value = 0xdead;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_b, &msg, 1000), -ENOMSG);
+  assert_int_equal(g_listener_b_count, 0);
+
+  msg.value = 7;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_b, &msg, 1000), 0);
+  assert_int_equal(g_listener_b_count, 1);
+}
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+/* Message subscriber: receives a copy of every message, in order, even if
+ * the channel is republished before the subscriber runs.
+ */
+
+static void test_msg_subscriber(FAR void **state)
+{
+  const struct zbus_channel *chan;
+  struct zbt_msg_s msg;
+  uint32_t expected[3] =
+  {
+    10, 20, 30
+  };
+
+  int i;
+
+  (void)state;
+  reset_all();
+
+  for (i = 0; i < 3; i++)
+    {
+      msg.seq = i;
+      msg.value = expected[i];
+      assert_int_equal(zbus_chan_pub(&zbt_chan_b, &msg, 1000), 0);
+    }
+
+  for (i = 0; i < 3; i++)
+    {
+      memset(&msg, 0, sizeof(msg));
+      assert_int_equal(zbus_sub_wait_msg(&zbt_msgsub_b, &chan, &msg, 1000),
+                       0);
+      assert_ptr_equal(chan, &zbt_chan_b);
+      assert_int_equal(msg.value, expected[i]);
+    }
+}
+#endif /* CONFIG_ZBUS_MSG_SUBSCRIBER */
+
+/* Float payload: float/double members must survive bit-exact through
+ * publish, the listener const view, read-back and the message subscriber
+ * copy; the validator exercises float math (non-finite rejection).
+ */
+
+static void test_float_payload(FAR void **state)
+{
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+  const struct zbus_channel *chan;
+#endif
+  struct zbt_imu_msg_s msg;
+  struct zbt_imu_msg_s rd;
+  int ret;
+
+  (void)state;
+
+  g_imu_listener_count = 0;
+
+  msg.accel[0] = 0.5f;
+  msg.accel[1] = -9.80665f;
+  msg.accel[2] = 3.1415927f;
+  msg.magnitude = 9.83180020299;
+  msg.seq = 1;
+
+  ret = zbus_chan_pub(&zbt_chan_imu, &msg, 1000);
+  assert_int_equal(ret, 0);
+
+  /* Listener saw a bit-exact copy */
+
+  assert_int_equal(g_imu_listener_count, 1);
+  assert_true(g_imu_listener_last.accel[0] == msg.accel[0]);
+  assert_true(g_imu_listener_last.accel[1] == msg.accel[1]);
+  assert_true(g_imu_listener_last.accel[2] == msg.accel[2]);
+  assert_true(g_imu_listener_last.magnitude == msg.magnitude);
+
+  /* Read-back from channel storage */
+
+  memset(&rd, 0, sizeof(rd));
+  ret = zbus_chan_read(&zbt_chan_imu, &rd, 500);
+  assert_int_equal(ret, 0);
+  assert_true(rd.accel[0] == msg.accel[0]);
+  assert_true(rd.accel[1] == msg.accel[1]);
+  assert_true(rd.accel[2] == msg.accel[2]);
+  assert_true(rd.magnitude == msg.magnitude);
+  assert_int_equal(rd.seq, 1);
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+  /* Message subscriber received a bit-exact copy through the mq */
+
+  memset(&rd, 0, sizeof(rd));
+  ret = zbus_sub_wait_msg(&zbt_imu_msgsub, &chan, &rd, 1000);
+  assert_int_equal(ret, 0);
+  assert_ptr_equal(chan, &zbt_chan_imu);
+  assert_true(rd.accel[0] == msg.accel[0]);
+  assert_true(rd.accel[1] == msg.accel[1]);
+  assert_true(rd.accel[2] == msg.accel[2]);
+  assert_true(rd.magnitude == msg.magnitude);
+#endif
+
+  /* Validator rejects non-finite samples with -ENOMSG */
+
+  msg.accel[1] = NAN;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_imu, &msg, 1000), -ENOMSG);
+  assert_int_equal(g_imu_listener_count, 1);
+
+  msg.accel[1] = INFINITY;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_imu, &msg, 1000), -ENOMSG);
+  assert_int_equal(g_imu_listener_count, 1);
+}
+
+/* Timer-driven publisher: a kernel timer interrupt wakes a sampling
+ * thread through a signal, which publishes float samples; the message
+ * subscriber receives every sample in order, bit-exact, and the listener
+ * sees each publication.
+ */
+
+static void test_timer_driven_publisher(FAR void **state)
+{
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+  const struct zbus_channel *chan;
+  struct zbt_imu_msg_s rd;
+  int i;
+#endif
+  pthread_t thread;
+
+  (void)state;
+
+  g_imu_listener_count = 0;
+  g_sampler_published = 0;
+  g_sampler_errors = 0;
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+  while (zbus_sub_wait_msg(&zbt_imu_msgsub, &chan, &rd, ZBUS_NO_WAIT) == 0)
+    {
+    }
+#endif
+
+  assert_int_equal(pthread_create(&thread, NULL, sampler_thread, NULL), 0);
+
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+  for (i = 1; i <= ZBT_SAMPLER_SAMPLES; i++)
+    {
+      memset(&rd, 0, sizeof(rd));
+      assert_int_equal(zbus_sub_wait_msg(&zbt_imu_msgsub, &chan, &rd, 1000),
+                       0);
+      assert_ptr_equal(chan, &zbt_chan_imu);
+      assert_int_equal(rd.seq, i);
+      assert_true(rd.accel[0] == 0.5f * (i - 1));
+      assert_true(rd.accel[1] == -9.5f);
+      assert_true(rd.magnitude == 9.5 + (i - 1));
+    }
+#endif
+
+  assert_int_equal(pthread_join(thread, NULL), 0);
+  assert_int_equal(g_sampler_errors, 0);
+  assert_int_equal(g_sampler_published, ZBT_SAMPLER_SAMPLES);
+  assert_int_equal(g_imu_listener_count, ZBT_SAMPLER_SAMPLES);
+  assert_int_equal(g_imu_listener_last.seq, ZBT_SAMPLER_SAMPLES);
+}
+
+/* Notification order: the observers listed in the channel definition are
+ * notified in the order they appear, and an observation bound out of line
+ * with ZBUS_CHAN_ADD_OBS() comes after all of them.
+ */
+
+static void test_observer_order(FAR void **state)
+{
+  struct zbt_msg_s msg;
+
+  (void)state;
+
+  g_order_seq = 0;
+  g_order_first = 0;
+  g_order_second = 0;
+  g_order_added = 0;
+
+  msg.seq = 1;
+  msg.value = 7;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_ord, &msg, 1000), 0);
+
+  assert_int_equal(g_order_first, 1);
+  assert_int_equal(g_order_second, 2);
+  assert_int_equal(g_order_added, 3);
+}
+
+/* Claim/finish: direct access to the message memory; notify dispatches
+ * without publishing.
+ */
+
+static void test_claim_finish_notify(FAR void **state)
+{
+  struct zbt_msg_s *direct;
+  struct zbt_msg_s msg;
+
+  (void)state;
+  reset_all();
+
+  assert_int_equal(zbus_chan_claim(&zbt_chan_a, 500), 0);
+
+  direct = zbus_chan_msg(&zbt_chan_a);
+  assert_non_null(direct);
+  direct->value = 55;
+
+  assert_int_equal(zbus_chan_finish(&zbt_chan_a), 0);
+
+  /* No notification happened yet */
+
+  assert_int_equal(g_listener_a_count, 0);
+
+  /* Force the notification: the listener must see value 55 */
+
+  assert_int_equal(zbus_chan_notify(&zbt_chan_a, 1000), 0);
+  assert_int_equal(g_listener_a_count, 1);
+  assert_int_equal(g_listener_a_last, 55);
+
+  memset(&msg, 0, sizeof(msg));
+  assert_int_equal(zbus_chan_read(&zbt_chan_a, &msg, 500), 0);
+  assert_int_equal(msg.value, 55);
+
+  drain_subscriber(&zbt_sub_a);
+}
+
+/* Notification masks: a masked observer is skipped; unrelated pairs
+ * return -ESRCH.
+ */
+
+static void test_masks(FAR void **state)
+{
+  struct zbt_msg_s msg;
+  bool masked;
+
+  (void)state;
+  reset_all();
+
+  assert_int_equal(zbus_obs_set_chan_notification_mask(&zbt_listener_a,
+                                                       &zbt_chan_a, true),
+                   0);
+  assert_int_equal(zbus_obs_is_chan_notification_masked(&zbt_listener_a,
+                                                        &zbt_chan_a,
+                                                        &masked), 0);
+  assert_true(masked);
+
+  msg.seq = 1;
+  msg.value = 77;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, 1000), 0);
+  assert_int_equal(g_listener_a_count, 0);
+
+  assert_int_equal(zbus_obs_set_chan_notification_mask(&zbt_listener_a,
+                                                       &zbt_chan_a, false),
+                   0);
+
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, 1000), 0);
+  assert_int_equal(g_listener_a_count, 1);
+
+  /* listener_b does not observe chan_a */
+
+  assert_int_equal(zbus_obs_set_chan_notification_mask(&zbt_listener_b,
+                                                       &zbt_chan_a, true),
+                   -ESRCH);
+
+  drain_subscriber(&zbt_sub_a);
+}
+
+/* Observer enable/disable */
+
+static void test_enable_disable(FAR void **state)
+{
+  struct zbt_msg_s msg;
+  bool enabled;
+
+  (void)state;
+  reset_all();
+
+  assert_int_equal(zbus_obs_set_enable(&zbt_listener_a, false), 0);
+  assert_int_equal(zbus_obs_is_enabled(&zbt_listener_a, &enabled), 0);
+  assert_false(enabled);
+
+  msg.seq = 1;
+  msg.value = 88;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, 1000), 0);
+  assert_int_equal(g_listener_a_count, 0);
+
+  assert_int_equal(zbus_obs_set_enable(&zbt_listener_a, true), 0);
+
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, 1000), 0);
+  assert_int_equal(g_listener_a_count, 1);
+
+  drain_subscriber(&zbt_sub_a);
+}
+
+#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
+/* Runtime observers: add/remove, duplicate detection */
+
+static void test_runtime_observers(FAR void **state)
+{
+  struct zbt_msg_s msg;
+
+  (void)state;
+  reset_all();
+
+  assert_int_equal(zbus_chan_add_obs(&zbt_chan_a, &zbt_listener_rt, 500),
+                   0);
+
+  /* Duplicates: already a runtime observer / already a static observer */
+
+  assert_int_equal(zbus_chan_add_obs(&zbt_chan_a, &zbt_listener_rt, 500),
+                   -EALREADY);
+  assert_int_equal(zbus_chan_add_obs(&zbt_chan_a, &zbt_listener_a, 500),
+                   -EEXIST);
+
+  msg.seq = 1;
+  msg.value = 99;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, 1000), 0);
+  assert_int_equal(g_listener_rt_count, 1);
+  assert_int_equal(g_listener_a_count, 1);
+
+  assert_int_equal(zbus_chan_rm_obs(&zbt_chan_a, &zbt_listener_rt, 500),
+                   0);
+
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, 1000), 0);
+  assert_int_equal(g_listener_rt_count, 1);
+
+  assert_int_equal(zbus_chan_rm_obs(&zbt_chan_a, &zbt_listener_rt, 500),
+                   -ENODATA);
+
+  drain_subscriber(&zbt_sub_a);
+}
+#endif /* CONFIG_ZBUS_RUNTIME_OBSERVERS */
+
+#ifdef CONFIG_ZBUS_ASYNC_LISTENER
+/* Async listener: callback runs on the listener's task with a copy of
+ * the message; a burst of publishes is delivered completely and in order.
+ */
+
+static void test_async_listener(FAR void **state)
+{
+  struct zbt_msg_s msg;
+  int i;
+
+  (void)state;
+  reset_all();
+
+  msg.seq = 1;
+  msg.value = 4242;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_b, &msg, 1000), 0);
+
+  for (i = 0; i < 100 && g_async_count < 1; i++)
+    {
+      usleep(10 * 1000);
+    }
+
+  assert_int_equal(g_async_count, 1);
+  assert_int_equal(g_async_last, 4242);
+
+  /* Burst: all copies must be delivered */
+
+  for (i = 1; i <= 3; i++)
+    {
+      msg.value = 4242 + i;
+      assert_int_equal(zbus_chan_pub(&zbt_chan_b, &msg, 1000), 0);
+    }
+
+  for (i = 0; i < 100 && g_async_count < 4; i++)
+    {
+      usleep(10 * 1000);
+    }
+
+  assert_int_equal(g_async_count, 4);
+  assert_int_equal(g_async_last, 4245);
+}
+#endif /* CONFIG_ZBUS_ASYNC_LISTENER */
+
+#ifdef CONFIG_ZBUS_CHANNEL_NAME
+/* Channel lookup by name */
+
+static void test_from_name(FAR void **state)
+{
+  (void)state;
+
+  assert_ptr_equal(zbus_chan_from_name("zbt_chan_a"), &zbt_chan_a);
+  assert_ptr_equal(zbus_chan_from_name("zbt_chan_b"), &zbt_chan_b);
+  assert_null(zbus_chan_from_name("does_not_exist"));
+
+  assert_string_equal(zbus_chan_name(&zbt_chan_a), "zbt_chan_a");
+}
+#endif
+
+#ifdef CONFIG_ZBUS_CHANNEL_ID
+/* Channel lookup by numeric id */
+
+static void test_from_id(FAR void **state)
+{
+  (void)state;
+
+  assert_ptr_equal(zbus_chan_from_id(42), &zbt_chan_c);
+  assert_null(zbus_chan_from_id(0xfffffff0));
+  assert_null(zbus_chan_from_id(ZBUS_CHAN_ID_INVALID));
+}
+#endif
+
+/* Iteration over channels and observers */
+
+static bool count_channel(const struct zbus_channel *chan, void *user_data)
+{
+  int *count = user_data;
+
+  (void)chan;
+  (*count)++;
+  return true;
+}
+
+static bool count_observer(const struct zbus_observer *obs, void *user_data)
+{
+  int *count = user_data;
+
+  (void)obs;
+  (*count)++;
+  return true;
+}
+
+static void test_iterate(FAR void **state)
+{
+  int channels = 0;
+  int observers = 0;
+
+  (void)state;
+
+  assert_true(zbus_iterate_over_channels_with_user_data(count_channel,
+                                                        &channels));
+  assert_true(zbus_iterate_over_observers_with_user_data(count_observer,
+                                                         &observers));
+
+  /* At least the three test channels and four test observers exist
+   * (other zbus users may add more to the image).
+   */
+
+  assert_true(channels >= 3);
+  assert_true(observers >= 4);
+}
+
+/* Message metadata accessors */
+
+static void test_accessors(FAR void **state)
+{
+  (void)state;
+
+  assert_int_equal(zbus_chan_msg_size(&zbt_chan_a),
+                   sizeof(struct zbt_msg_s));
+  assert_ptr_equal(zbus_chan_user_data(&zbt_chan_b), &g_user_word);
+  assert_null(zbus_chan_user_data(&zbt_chan_a));
+}
+
+/* Timeout semantics: subscriber queue overflow reports -ENOMSG on
+ * no-wait publish; empty queue reports -ENOMSG (no-wait) or -EAGAIN
+ * (timed out).
+ */
+
+static void test_timeouts(FAR void **state)
+{
+  const struct zbus_channel *chan;
+  struct zbt_msg_s msg;
+  int i;
+
+  (void)state;
+  reset_all();
+
+  /* zbt_sub_a queue depth is 4: four publishes succeed... */
+
+  msg.seq = 1;
+  for (i = 0; i < 4; i++)
+    {
+      msg.value = i;
+      assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, ZBUS_NO_WAIT), 0);
+    }
+
+  /* ...the fifth overflows the subscriber queue.  The publish itself
+   * happens (the message is copied and the listener notified); the
+   * -ENOMSG is the collected delivery error, matching the Zephyr VDED
+   * semantics.
+   */
+
+  msg.value = 4;
+  assert_int_equal(zbus_chan_pub(&zbt_chan_a, &msg, ZBUS_NO_WAIT), -ENOMSG);
+  assert_int_equal(g_listener_a_count, 5);
+
+  /* Drain the four queued notifications */
+
+  for (i = 0; i < 4; i++)
+    {
+      assert_int_equal(zbus_sub_wait(&zbt_sub_a, &chan, ZBUS_NO_WAIT), 0);
+      assert_ptr_equal(chan, &zbt_chan_a);
+    }
+
+  /* Empty queue: no-wait -> -ENOMSG, timed -> -EAGAIN */
+
+  assert_int_equal(zbus_sub_wait(&zbt_sub_a, &chan, ZBUS_NO_WAIT), -ENOMSG);
+  assert_int_equal(zbus_sub_wait(&zbt_sub_a, &chan, 50), -EAGAIN);
+}
+
+/****************************************************************************
+ * Public Functions
+ ****************************************************************************/
+
+int main(int argc, FAR char *argv[])
+{
+  const struct CMUnitTest tests[] =
+  {
+    cmocka_unit_test(test_pub_read_listener_subscriber),
+    cmocka_unit_test(test_multi_channel_isolation),
+    cmocka_unit_test(test_validator),
+#ifdef CONFIG_ZBUS_MSG_SUBSCRIBER
+    cmocka_unit_test(test_msg_subscriber),
+#endif
+    cmocka_unit_test(test_float_payload),
+    cmocka_unit_test(test_timer_driven_publisher),
+    cmocka_unit_test(test_observer_order),
+    cmocka_unit_test(test_claim_finish_notify),
+    cmocka_unit_test(test_masks),
+    cmocka_unit_test(test_enable_disable),
+#ifdef CONFIG_ZBUS_RUNTIME_OBSERVERS
+    cmocka_unit_test(test_runtime_observers),
+#endif
+#ifdef CONFIG_ZBUS_ASYNC_LISTENER
+    cmocka_unit_test(test_async_listener),
+#endif
+#ifdef CONFIG_ZBUS_CHANNEL_NAME
+    cmocka_unit_test(test_from_name),
+#endif
+#ifdef CONFIG_ZBUS_CHANNEL_ID
+    cmocka_unit_test(test_from_id),
+#endif
+    cmocka_unit_test(test_iterate),
+    cmocka_unit_test(test_accessors),
+    cmocka_unit_test(test_timeouts),
+  };
+
+  return cmocka_run_group_tests(tests, NULL, NULL);
+}

Reply via email to