libxr  1.0
Want to be the best embedded framework
Loading...
Searching...
No Matches
topic.hpp
1#pragma once
2
3#include <atomic>
4#include <concepts>
5#include <cstddef>
6#include <cstdint>
7
8#include "libxr_cb.hpp"
9#include "libxr_def.hpp"
10#include "libxr_time.hpp"
11#include "libxr_type.hpp"
12#include "lockfree_list.hpp"
13#include "mutex.hpp"
14#include "queue.hpp"
15#include "rbt.hpp"
16#include "semaphore.hpp"
17#include "thread.hpp"
18
19namespace LibXR
20{
26template <typename Data>
27concept TopicPayload =
28 !std::is_reference_v<Data> && !std::is_const_v<Data> && !std::is_volatile_v<Data> &&
29 std::is_object_v<Data> && std::is_default_constructible_v<Data> &&
30 std::is_copy_assignable_v<Data> && std::is_trivially_destructible_v<Data>;
31
37template <typename Data>
38constexpr void CheckTopicPayload()
39{
40 static_assert(TopicPayload<Data>,
41 "LibXR::Topic typed payload must be a non-cv/ref object type that is "
42 "default-constructible, copy-assignable, and trivially destructible.");
43}
44
55class Topic
56{
61 enum class LockState : uint32_t
62 {
63 UNLOCKED = 0,
64 LOCKED = 1,
66 USE_MUTEX = UINT32_MAX
68 };
69
70 public:
81 struct Block
82 {
83 std::atomic<LockState>
86 TypeID::ID
88 uint32_t payload_size;
92 uint32_t crc32;
94 };
95
96#ifndef __DOXYGEN__
101 struct PackedDataHeader;
102
109 template <typename Data>
110 class PackedData;
111 static constexpr uint8_t PACKET_PREFIX =
112 0x5A;
113 static constexpr uint8_t PACKET_VERSION =
114 0x01;
115 static constexpr size_t PACK_BASE_SIZE =
116 17;
118#endif
119
126
133 template <typename Data>
135 {
136 static_assert(TopicPayload<Data>);
138 Data* data;
140 };
141
159
166 template <typename Data>
167 struct Message
168 {
169 static_assert(TopicPayload<Data>);
171 Data data;
172 };
173
179 static void Lock(TopicHandle topic);
180
186 static void Unlock(TopicHandle topic);
187
193 static void LockFromCallback(TopicHandle topic);
194
200 static void UnlockFromCallback(TopicHandle topic);
201
206 class Domain
207 {
208 public:
213 Domain(const char* name);
214
218 };
219
224 enum class SuberType : uint8_t
225 {
226 SYNC,
227 ASYNC,
228 QUEUE,
229 CALLBACK,
230 };
231
237 {
239 };
240
246 struct SyncBlock;
247
254 template <typename Data>
256
261 enum class ASyncSubscriberState : uint32_t;
262
268 struct ASyncBlock;
269
276 template <typename Data>
278
284 struct QueueBlock;
285
290 class QueuedSubscriber;
291
297 class Callback;
298
304 struct CallbackBlock;
305
311 class Server;
312
317 void RegisterCallback(Callback& cb);
318
324
328 Topic();
329
343 Topic(const char* name, TypeID::ID payload_type_id, size_t payload_size,
344 size_t payload_alignment, Domain* domain = nullptr, bool multi_publisher = false);
345
356 template <typename Data>
357 static Topic CreateTopic(const char* name, Domain* domain = nullptr,
358 bool multi_publisher = false)
359 {
361 return Topic(name, TypeID::GetID<Data>(), sizeof(Data), alignof(Data), domain,
362 multi_publisher);
363 }
364
370 Topic(TopicHandle topic);
371
378 static TopicHandle Find(const char* name, Domain* domain = nullptr);
379
390 template <typename Data>
391 static TopicHandle FindOrCreate(const char* name, Domain* domain = nullptr,
392 bool multi_publisher = false)
393 {
395 auto topic = Find(name, domain);
396 if (topic != nullptr)
397 {
399 if (multi_publisher && !topic->data_.mutex)
400 {
401 ASSERT(false);
402 }
403 }
404 else
405 {
406 topic = CreateTopic<Data>(name, domain, multi_publisher).block_;
407 }
408 return topic;
409 }
410
415 [[nodiscard]] size_t PayloadSize() const
416 {
417 ASSERT(block_ != nullptr);
418 return block_->data_.payload_size;
419 }
420
426 [[nodiscard]] size_t PayloadAlignment() const
427 {
428 ASSERT(block_ != nullptr);
429 return block_->data_.payload_alignment;
430 }
431
438 template <typename Data>
439 void Publish(Data& data)
440 {
441 PublishTyped(data, NowTimestamp(), false, false);
442 }
443
451 template <typename Data>
452 void Publish(Data& data, MicrosecondTimestamp timestamp)
453 {
454 PublishTyped(data, timestamp, false, false);
455 }
456
464 template <typename Data>
465 void PublishFromCallback(Data& data, bool in_isr)
466 {
467 PublishTyped(data, NowTimestamp(), true, in_isr);
468 }
469
478 template <typename Data>
479 void PublishFromCallback(Data& data, MicrosecondTimestamp timestamp, bool in_isr)
480 {
481 PublishTyped(data, timestamp, true, in_isr);
482 }
483
493 void PublishBytesFromServer(void* payload_addr, size_t payload_size,
494 MicrosecondTimestamp timestamp)
495 {
496 PublishServerBytes(payload_addr, payload_size, timestamp, false, false);
497 }
498
510 void PublishBytesFromServerCallback(void* payload_addr, size_t payload_size,
511 MicrosecondTimestamp timestamp, bool in_isr)
512 {
513 PublishServerBytes(payload_addr, payload_size, timestamp, true, in_isr);
514 }
515
524 template <typename Data>
525 ErrorCode PackData(const Data& data, PackedData<Data>& packet)
526 {
527 return PackData(data, packet, NowTimestamp());
528 }
529
539 template <typename Data>
540 ErrorCode PackData(const Data& data, PackedData<Data>& packet,
541 MicrosecondTimestamp timestamp);
542
552 {
553 return PackRaw(data, packet, NowTimestamp());
554 }
555
566 {
567 if (block_ == nullptr || data.addr_ == nullptr || packet.addr_ == nullptr)
568 {
569 return ErrorCode::PTR_NULL;
570 }
571
572 if (data.size_ != block_->data_.payload_size)
573 {
574 return ErrorCode::SIZE_ERR;
575 }
576
577 if (packet.size_ < PACK_BASE_SIZE + data.size_)
578 {
579 return ErrorCode::NO_BUFF;
580 }
581
582 PackBytes(block_->data_.crc32, packet, timestamp, data);
583 return ErrorCode::OK;
584 }
585
593 static TopicHandle WaitTopic(const char* name, uint32_t timeout = UINT32_MAX,
594 Domain* domain = nullptr);
595
601 operator TopicHandle() { return block_; }
602
607 uint32_t GetKey() const;
608
609 private:
610 TopicHandle block_ = nullptr;
612
613 static inline RBTree<uint32_t>* domain_ =
614 nullptr;
615 static inline Domain* def_domain_ = nullptr;
616
620 static void EnsureDomainRegistry();
621
626 static Domain* EnsureDefaultDomain();
627
635 static void CheckServerPublishContract(TopicHandle topic, void* payload_addr,
636 size_t payload_size)
637 {
638 ASSERT(topic != nullptr);
639 ASSERT(payload_addr != nullptr);
640 ASSERT(payload_size == topic->data_.payload_size);
641 ASSERT(topic->data_.payload_alignment != 0);
642 ASSERT(reinterpret_cast<uintptr_t>(payload_addr) % topic->data_.payload_alignment ==
643 0);
644 }
645
652 template <typename Data>
653 static void CheckSubscriberType(Topic topic)
654 {
656 ASSERT(topic.block_ != nullptr);
657 ASSERT(topic.block_->data_.payload_type_id == TypeID::GetID<Data>());
658 ASSERT(topic.block_->data_.payload_size == sizeof(Data));
659 ASSERT(topic.block_->data_.payload_alignment == alignof(Data));
660 }
661
668 template <typename Data>
670 {
672 return new Data;
673 }
674
683 template <typename Data>
684 static void CopyPayload(void* dst, void* payload_addr)
685 {
687 ASSERT(dst != nullptr);
688 ASSERT(payload_addr != nullptr);
689 *reinterpret_cast<Data*>(dst) = *reinterpret_cast<Data*>(payload_addr);
690 }
691
702 template <typename Data>
703 void PublishTyped(Data& data, MicrosecondTimestamp timestamp, bool from_callback,
704 bool in_isr)
705 {
707
708 if (from_callback)
709 {
711 }
712 else
713 {
714 Lock(block_);
715 }
716
717 CheckPublishContract(block_, TypeID::GetID<Data>(), sizeof(Data), alignof(Data));
718 DispatchSubscribers(block_, timestamp, &data, from_callback, in_isr);
719
720 if (from_callback)
721 {
723 }
724 else
725 {
726 Unlock(block_);
727 }
728 }
729
741 static void CheckPublishContract(TopicHandle topic, TypeID::ID payload_type_id,
742 size_t payload_size, size_t payload_alignment);
743
752 static void PackBytes(uint32_t topic_name_crc32, RawData buffer,
753 MicrosecondTimestamp timestamp, ConstRawData data);
754
765 static void DispatchSubscriber(SuberBlock& block, MicrosecondTimestamp timestamp,
766 void* payload_addr, bool from_callback, bool in_isr);
767
779 static void DispatchSubscribers(TopicHandle topic, MicrosecondTimestamp timestamp,
780 void* payload_addr, bool from_callback, bool in_isr);
781
792 void PublishServerBytes(void* payload_addr, size_t payload_size,
793 MicrosecondTimestamp timestamp, bool from_callback, bool in_isr)
794 {
795 CheckServerPublishContract(block_, payload_addr, payload_size);
796
797 if (from_callback)
798 {
800 }
801 else
802 {
803 Lock(block_);
804 }
805
806 DispatchSubscribers(block_, timestamp, payload_addr, from_callback, in_isr);
807
808 if (from_callback)
809 {
811 }
812 else
813 {
814 Unlock(block_);
815 }
816 }
817};
818} // namespace LibXR
819
820#include "packet/packet.hpp"
821#include "server/server.hpp"
822#include "subscriber/async.hpp"
823#include "subscriber/callback.hpp"
824#include "subscriber/queue.hpp"
825#include "subscriber/sync.hpp"
只读原始数据视图 / Immutable raw data view
size_t size_
数据字节数 / Data size in bytes
const void * addr_
数据起始地址 / Data start address
链表实现,用于存储和管理数据节点。 A linked list implementation for storing and managing data nodes.
微秒时间戳 / Microsecond timestamp
互斥锁类,提供线程同步机制 (Mutex class providing thread synchronization mechanisms).
Definition mutex.hpp:18
红黑树的泛型数据节点,继承自 BaseNode (Generic data node for Red-Black Tree, inheriting from BaseNode).
Definition rbt.hpp:63
Data data_
存储的数据 (Stored data).
Definition rbt.hpp:98
红黑树实现,支持泛型键和值,并提供线程安全操作 (Red-Black Tree implementation supporting generic keys and values with thread...
Definition rbt.hpp:23
可写原始数据视图 / Mutable raw data view
size_t size_
数据字节数 / Data size in bytes
void * addr_
数据起始地址 / Data start address
先 StartWaiting(),再自己来取数据的订阅者 / Subscriber that first calls StartWaiting() and later pulls the data it...
Definition topic.hpp:277
每次发布时直接执行函数的订阅句柄 / Subscription handle that runs a function on each publish
Definition callback.hpp:16
topic 所属的命名域 / Naming domain that groups topics
Definition topic.hpp:207
Domain(const char *name)
构造一个 topic 域 / Construct one topic domain
Definition topic.cpp:92
RBTree< uint32_t >::Node< RBTree< uint32_t > > * node_
Definition topic.hpp:216
每次发布都往队列里塞一份数据的订阅者 / Subscriber that pushes one entry into a queue on each publish
Definition queue.hpp:25
将字节流解析成 packet 并发布到已注册 topic 的状态机 / State machine that parses byte streams into packets and publishes...
Definition server.hpp:27
调用 Wait() 收消息的订阅者 / Subscriber that receives messages by calling Wait()
Definition topic.hpp:255
发布订阅主题 / Publish-subscribe topic
Definition topic.hpp:56
SuberType
topic 支持的订阅者种类 / Subscriber kinds supported by a topic
Definition topic.hpp:225
@ SYNC
同步等待型订阅者。Synchronous wait-based subscriber.
@ ASYNC
异步本地缓冲型订阅者。Asynchronous local-buffer subscriber.
@ QUEUE
队列转发型订阅者。Queue-forwarding subscriber.
@ CALLBACK
回调执行型订阅者。Callback-executing subscriber.
static void PackBytes(uint32_t topic_name_crc32, RawData buffer, MicrosecondTimestamp timestamp, ConstRawData data)
将一段 payload 字节和 topic 元数据拼成 packet / Pack one payload byte range together with topic metadata into on...
Definition packet.cpp:42
void Publish(Data &data)
在普通上下文里发布一条消息,并自动取当前时间戳 / Publish one message in normal context and stamp it with the current time
Definition topic.hpp:439
void PublishFromCallback(Data &data, MicrosecondTimestamp timestamp, bool in_isr)
在回调或 ISR 路径里按指定时间戳发布一条消息 / Publish one message from callback or ISR context with an explicit timestam...
Definition topic.hpp:479
static Domain * def_domain_
缺省 topic 域。Default topic domain.
Definition topic.hpp:615
static void Unlock(TopicHandle topic)
在普通上下文里释放一个 topic 发布路径 / Unlock one topic publish path in normal context
Definition topic.cpp:50
static TopicHandle FindOrCreate(const char *name, Domain *domain=nullptr, bool multi_publisher=false)
按精确类型查找或创建一个 topic / Find or create one topic with an exact payload type
Definition topic.hpp:391
static void * AllocateSubscriberBuffer()
为订阅者分配一个长期存在的本地接收对象 / Allocate one long-lived local receive object for a subscriber
Definition topic.hpp:669
void PublishServerBytes(void *payload_addr, size_t payload_size, MicrosecondTimestamp timestamp, bool from_callback, bool in_isr)
PublishBytesFromServer*() 的共享实现 / Shared implementation behind PublishBytesFromServer*()
Definition topic.hpp:792
RBTree< uint32_t >::Node< Block > * TopicHandle
指向一个 topic 运行时状态块的句柄 / Handle pointing to one topic runtime state block
Definition topic.hpp:125
size_t PayloadSize() const
获取该 topic 固定 payload 字节数 / Get the fixed payload size of this topic
Definition topic.hpp:415
ASyncSubscriberState
异步订阅者本地缓冲区的状态 / State of the async subscriber's local buffer
Definition async.hpp:12
void PublishTyped(Data &data, MicrosecondTimestamp timestamp, bool from_callback, bool in_isr)
强类型发布入口的共享实现 / Shared implementation of typed publish entry points
Definition topic.hpp:703
ErrorCode PackRaw(ConstRawData data, RawData packet)
将一段 raw payload 按当前 topic 元数据打包成 packet / Pack a raw payload byte view into one packet using the curr...
Definition topic.hpp:551
static Domain * EnsureDefaultDomain()
确保默认域已创建 / Ensure the default domain exists
Definition topic.cpp:22
static MicrosecondTimestamp NowTimestamp()
读取当前时间戳 / Read the current timestamp
Definition publish.cpp:80
static Topic CreateTopic(const char *name, Domain *domain=nullptr, bool multi_publisher=false)
用精确类型创建或查找一个 topic / Create or look up one topic using one exact payload type
Definition topic.hpp:357
static void CopyPayload(void *dst, void *payload_addr)
按精确类型把一份 payload 拷到订阅者缓冲区 / Copy one payload into a subscriber buffer using the exact type
Definition topic.hpp:684
static void CheckSubscriberType(Topic topic)
断言订阅者看到的精确 payload 类型与 topic 契约一致 / Assert that the exact payload type seen by a subscriber matches t...
Definition topic.hpp:653
static void Lock(TopicHandle topic)
在普通上下文里锁住一个 topic 发布路径 / Lock one topic publish path in normal context
Definition topic.cpp:32
static void UnlockFromCallback(TopicHandle topic)
在回调或 ISR 路径里释放一个 topic 发布路径 / Unlock one topic publish path from callback or ISR context
Definition topic.cpp:80
static void LockFromCallback(TopicHandle topic)
在回调或 ISR 路径里锁住一个 topic 发布路径 / Lock one topic publish path from callback or ISR context
Definition topic.cpp:62
void PublishBytesFromServerCallback(void *payload_addr, size_t payload_size, MicrosecondTimestamp timestamp, bool in_isr)
供回调/ISR 上下文里的 packet/server 路径按字节发布一条消息 / Publish one packet/server message from callback/ISR context...
Definition topic.hpp:510
static TopicHandle WaitTopic(const char *name, uint32_t timeout=UINT32_MAX, Domain *domain=nullptr)
等待指定名称的 topic 出现 / Wait until a topic with the given name exists
Definition topic.cpp:190
size_t PayloadAlignment() const
获取该 topic payload 的对齐要求 / Get the payload alignment requirement of this topic
Definition topic.hpp:426
static void CheckServerPublishContract(TopicHandle topic, void *payload_addr, size_t payload_size)
校验 server 侧字节发布前提 / Check the preconditions of one server-side byte publish
Definition topic.hpp:635
static void CheckPublishContract(TopicHandle topic, TypeID::ID payload_type_id, size_t payload_size, size_t payload_alignment)
校验一次强类型发布的运行时契约 / Check the runtime contract of one typed publish
Definition publish.cpp:82
void PublishFromCallback(Data &data, bool in_isr)
在回调或 ISR 路径里发布一条消息,并自动取当前时间戳 / Publish one message from callback or ISR context and stamp it with the...
Definition topic.hpp:465
static void DispatchSubscriber(SuberBlock &block, MicrosecondTimestamp timestamp, void *payload_addr, bool from_callback, bool in_isr)
将一条消息分发给一个订阅块 / Dispatch one message to one subscriber block
Definition publish.cpp:12
uint32_t GetKey() const
读取 topic 键值 / Read the key value of this topic
Definition topic.cpp:211
TopicHandle block_
Definition topic.hpp:610
LockState
topic 发布路径的内部锁状态 / Internal lock state of the topic publish path
Definition topic.hpp:62
@ UNLOCKED
当前未持有发布锁。The publish path is currently unlocked.
void RegisterCallback(Callback &cb)
注册一个回调订阅者 / Register one callback subscriber
Definition callback.hpp:606
void Publish(Data &data, MicrosecondTimestamp timestamp)
在普通上下文里按指定时间戳发布一条消息 / Publish one message in normal context with an explicit timestamp
Definition topic.hpp:452
static TopicHandle Find(const char *name, Domain *domain=nullptr)
按名称查找一个已存在 topic / Find one existing topic by name
Definition topic.cpp:172
static RBTree< uint32_t > * domain_
全局 topic 域注册表。Global registry of topic domains.
Definition topic.hpp:613
static void DispatchSubscribers(TopicHandle topic, MicrosecondTimestamp timestamp, void *payload_addr, bool from_callback, bool in_isr)
将一条消息分发给一个 topic 上的全部订阅者 / Dispatch one message to all subscribers attached to one topic
Definition publish.cpp:69
Topic()
构造一个空 topic 视图 / Construct one empty topic view
Definition topic.cpp:114
void PublishBytesFromServer(void *payload_addr, size_t payload_size, MicrosecondTimestamp timestamp)
供 packet/server 路径按字节发布一条消息 / Publish one message from the packet/server path using bytes already arr...
Definition topic.hpp:493
ErrorCode PackRaw(ConstRawData data, RawData packet, MicrosecondTimestamp timestamp)
将一段 raw payload 按当前 topic 元数据和指定时间戳打包成 packet / Pack a raw payload byte view into one packet using th...
Definition topic.hpp:565
ErrorCode PackData(const Data &data, PackedData< Data > &packet)
将一个精确类型消息打包成 packet / Pack one exact-typed message into one packet using the topic's runtime contract...
Definition topic.hpp:525
static void EnsureDomainRegistry()
确保全局域注册表已创建 / Ensure the global domain registry exists
Definition topic.cpp:13
static ID GetID()
获取类型的唯一标识符 / Get a unique identifier for type T
topic 可承载 payload 的类型约束 / Type constraint for payloads carried by one topic
Definition topic.hpp:27
LibXR 命名空间
Definition ch32_can.hpp:14
ErrorCode
定义错误码枚举
@ SIZE_ERR
尺寸错误 | Size error
@ PTR_NULL
空指针 | Null pointer
@ NO_BUFF
缓冲区不足 | Insufficient buffer
@ OK
操作成功 | Operation successful
constexpr void CheckTopicPayload()
在模板上下文里断言 payload 类型满足 topic 契约 / Assert in template context that one payload type satisfies the topi...
Definition topic.hpp:38
异步订阅者自己挂的数据块 / Data block owned by one asynchronous subscriber
Definition async.hpp:26
topic 运行时状态块 / Runtime state block of one topic
Definition topic.hpp:82
std::atomic< LockState > busy
发布路径串行化状态。Publish-path serialization state.
Definition topic.hpp:84
Mutex * mutex
多发布者主题使用的互斥量。Mutex used by multi-publisher topics.
Definition topic.hpp:93
TypeID::ID payload_type_id
精确 payload 类型标识。Exact payload type identifier.
Definition topic.hpp:87
uint32_t payload_alignment
Definition topic.hpp:90
uint32_t payload_size
Definition topic.hpp:88
uint32_t crc32
主题名 CRC32 键。CRC32 key of the topic name.
Definition topic.hpp:92
LockFreeList subers
已挂接订阅者链表。List of attached subscribers.
Definition topic.hpp:85
挂在 topic 订阅链表里的回调记录 / Callback record stored in the topic subscriber list
Definition callback.hpp:566
带时间戳和 payload 副本的消息对象 / Message object carrying a timestamp and a payload copy
Definition topic.hpp:168
Data data
payload 对象副本。Copied payload object.
Definition topic.hpp:171
MicrosecondTimestamp timestamp
消息时间戳。Message timestamp.
Definition topic.hpp:170
带时间戳和 payload 指针的只读消息视图 / Read-only message view carrying a timestamp and a payload pointer
Definition topic.hpp:135
MicrosecondTimestamp timestamp
消息时间戳。Message timestamp.
Definition topic.hpp:137
队列订阅者自己挂的数据块 / Data block owned by one queued subscriber
Definition queue.hpp:12
带时间戳和只读原始 payload 视图的消息 / Message carrying a timestamp and a read-only raw payload view
Definition topic.hpp:154
MicrosecondTimestamp timestamp
消息时间戳。Message timestamp.
Definition topic.hpp:155
所有订阅块共用的公共头 / Common header shared by all subscriber blocks
Definition topic.hpp:237
SuberType type
订阅块的具体种类。Concrete kind of this subscriber block.
Definition topic.hpp:238
同步订阅者自己挂的数据块 / Data block owned by one synchronous subscriber
Definition sync.hpp:12
队列模块聚合入口。