3#if defined(LIBXR_SYSTEM_POSIX_HOST)
6#include <linux/futex.h>
10#include <sys/syscall.h>
29#include "libxr_def.hpp"
38enum class LinuxSharedSubscriberMode : uint8_t
50struct LinuxSharedTopicConfig
52 uint32_t slot_num = 64;
53 uint32_t subscriber_num = 8;
73template <
typename TopicData>
74class LinuxSharedTopic :
public Topic
76 static_assert(std::is_trivially_copyable<TopicData>::value,
77 "LinuxSharedTopic requires trivially copyable data");
78 static_assert(std::atomic<uint32_t>::is_always_lock_free,
79 "LinuxSharedTopic requires lock-free 32-bit atomics");
80 static_assert(std::atomic<uint64_t>::is_always_lock_free,
81 "LinuxSharedTopic requires lock-free 64-bit atomics");
83 enum class SharedDataState : uint8_t
92 using Data = SharedData;
93 static constexpr const char* DEFAULT_DOMAIN_NAME =
"libxr_def_domain";
104 using Data = SharedData;
110 Subscriber() =
default;
120 explicit Subscriber(
const char* name, LinuxSharedSubscriberMode mode =
121 LinuxSharedSubscriberMode::BROADCAST_FULL)
122 : owned_topic_(new LinuxSharedTopic(name))
124 if (Attach(*owned_topic_, mode) != ErrorCode::OK)
127 owned_topic_ =
nullptr;
131 Subscriber(
const char* name,
const char* domain_name,
132 LinuxSharedSubscriberMode mode = LinuxSharedSubscriberMode::BROADCAST_FULL)
133 : owned_topic_(new LinuxSharedTopic(name, domain_name))
135 if (Attach(*owned_topic_, mode) != ErrorCode::OK)
138 owned_topic_ =
nullptr;
142 Subscriber(
const char* name, Topic::Domain& domain,
143 LinuxSharedSubscriberMode mode = LinuxSharedSubscriberMode::BROADCAST_FULL)
144 : owned_topic_(new LinuxSharedTopic(name, domain))
146 if (Attach(*owned_topic_, mode) != ErrorCode::OK)
149 owned_topic_ =
nullptr;
160 LinuxSharedTopic& topic,
161 LinuxSharedSubscriberMode mode = LinuxSharedSubscriberMode::BROADCAST_FULL)
163 (void)Attach(topic, mode);
169 ~Subscriber() { Reset(); }
171 Subscriber(
const Subscriber&) =
delete;
172 Subscriber& operator=(
const Subscriber&) =
delete;
174 Subscriber(Subscriber&& other)
noexcept { *
this = std::move(other); }
176 Subscriber& operator=(Subscriber&& other)
noexcept
185 topic_ = other.topic_;
186 owned_topic_ = other.owned_topic_;
187 subscriber_index_ = other.subscriber_index_;
188 current_slot_index_ = other.current_slot_index_;
189 current_sequence_ = other.current_sequence_;
190 current_timestamp_ = other.current_timestamp_;
192 other.topic_ =
nullptr;
193 other.owned_topic_ =
nullptr;
194 other.subscriber_index_ = INVALID_INDEX;
195 other.current_slot_index_ = INVALID_INDEX;
196 other.current_sequence_ = 0;
197 other.current_timestamp_ = MicrosecondTimestamp();
206 bool Valid()
const {
return topic_ !=
nullptr && subscriber_index_ != INVALID_INDEX; }
215 ErrorCode Wait(uint32_t timeout_ms = UINT32_MAX)
219 return ErrorCode::STATE_ERR;
222 const uint64_t deadline_ms =
223 (timeout_ms == UINT32_MAX) ? 0 : (NowMonotonicMs() + timeout_ms);
225 Descriptor desc = {};
228 ErrorCode pop_ans = topic_->TryPopDescriptor(subscriber_index_, desc);
229 if (pop_ans == ErrorCode::OK)
232 topic_->HoldSlot(subscriber_index_, desc.slot_index);
233 current_slot_index_ = desc.slot_index;
234 current_sequence_ = desc.sequence;
235 current_timestamp_ = topic_->SlotTimestamp(desc.slot_index);
236 return ErrorCode::OK;
239 uint32_t wait_ms = UINT32_MAX;
240 if (timeout_ms != UINT32_MAX)
242 const uint64_t now_ms = NowMonotonicMs();
243 if (now_ms >= deadline_ms)
245 return ErrorCode::TIMEOUT;
247 wait_ms =
static_cast<uint32_t
>(deadline_ms - now_ms);
251 topic_->WaitReady(topic_->subscribers_[subscriber_index_], wait_ms);
252 if (wait_ans == ErrorCode::OK)
268 ErrorCode Wait(SharedData& data, uint32_t timeout_ms = UINT32_MAX)
274 return ErrorCode::STATE_ERR;
277 const uint64_t deadline_ms =
278 (timeout_ms == UINT32_MAX) ? 0 : (NowMonotonicMs() + timeout_ms);
280 Descriptor desc = {};
283 ErrorCode pop_ans = topic_->TryPopDescriptor(subscriber_index_, desc);
284 if (pop_ans == ErrorCode::OK)
286 data.topic_ = topic_;
287 data.slot_index_ = desc.slot_index;
288 data.sequence_ = desc.sequence;
289 data.state_ = SharedDataState::SUBSCRIBER;
290 data.subscriber_index_ = subscriber_index_;
291 topic_->HoldSlot(subscriber_index_, desc.slot_index);
292 return ErrorCode::OK;
295 uint32_t wait_ms = UINT32_MAX;
296 if (timeout_ms != UINT32_MAX)
298 const uint64_t now_ms = NowMonotonicMs();
299 if (now_ms >= deadline_ms)
301 return ErrorCode::TIMEOUT;
303 wait_ms =
static_cast<uint32_t
>(deadline_ms - now_ms);
307 topic_->WaitReady(topic_->subscribers_[subscriber_index_], wait_ms);
308 if (wait_ans == ErrorCode::OK)
321 TopicData* GetData()
const
323 if (!Valid() || current_slot_index_ == INVALID_INDEX)
328 return &topic_->payloads_[current_slot_index_];
334 uint64_t GetSequence()
const {
return current_sequence_; }
339 MicrosecondTimestamp GetTimestamp()
const {
return current_timestamp_; }
344 uint32_t GetPendingNum()
const
351 const SubscriberControl& control = topic_->subscribers_[subscriber_index_];
352 const uint32_t head = control.queue_head.load(std::memory_order_acquire);
353 const uint32_t tail = control.queue_tail.load(std::memory_order_acquire);
358 return topic_->queue_capacity_ - (head - tail);
365 uint64_t GetDropNum()
const
372 return topic_->subscribers_[subscriber_index_].dropped_messages.load(
373 std::memory_order_acquire);
381 if (!Valid() || current_slot_index_ == INVALID_INDEX)
386 topic_->ClearHeldSlot(subscriber_index_, current_slot_index_);
387 topic_->ReleaseSlot(current_slot_index_);
388 current_slot_index_ = INVALID_INDEX;
389 current_sequence_ = 0;
390 current_timestamp_ = MicrosecondTimestamp();
404 topic_->UnregisterBalancedSubscriber(subscriber_index_);
405 topic_->subscribers_[subscriber_index_].active.store(0, std::memory_order_release);
406 topic_->subscribers_[subscriber_index_].owner_pid.store(0,
407 std::memory_order_release);
408 topic_->subscribers_[subscriber_index_].owner_starttime.store(
409 0, std::memory_order_release);
411 Descriptor desc = {};
412 while (topic_->TryPopDescriptor(subscriber_index_, desc) == ErrorCode::OK)
414 topic_->ReleaseSlot(desc.slot_index);
421 owned_topic_ =
nullptr;
422 subscriber_index_ = INVALID_INDEX;
423 current_slot_index_ = INVALID_INDEX;
424 current_sequence_ = 0;
425 current_timestamp_ = MicrosecondTimestamp();
429 ErrorCode Attach(LinuxSharedTopic& topic, LinuxSharedSubscriberMode mode)
435 return ErrorCode::STATE_ERR;
438 if (topic.self_identity_.starttime == 0)
440 return ErrorCode::STATE_ERR;
443 for (uint32_t i = 0; i < topic.subscriber_capacity_; ++i)
445 uint32_t expected = 0;
446 auto& active = topic.subscribers_[i].active;
447 if (active.compare_exchange_strong(expected, 1, std::memory_order_acq_rel,
448 std::memory_order_relaxed))
450 topic.subscribers_[i].queue_head.store(0, std::memory_order_release);
451 topic.subscribers_[i].queue_tail.store(0, std::memory_order_release);
452 topic.subscribers_[i].ready_sem_count.store(0, std::memory_order_release);
453 topic.subscribers_[i].dropped_messages.store(0, std::memory_order_release);
454 topic.subscribers_[i].owner_pid.store(topic.self_identity_.pid,
455 std::memory_order_release);
456 topic.subscribers_[i].owner_starttime.store(topic.self_identity_.starttime,
457 std::memory_order_release);
458 topic.subscribers_[i].held_slot.store(INVALID_INDEX, std::memory_order_release);
459 topic.subscribers_[i].mode.store(
static_cast<uint32_t
>(mode),
460 std::memory_order_release);
461 if (mode == LinuxSharedSubscriberMode::BALANCE_RR)
463 const ErrorCode join_ans = topic.RegisterBalancedSubscriber(i);
464 if (join_ans != ErrorCode::OK)
466 topic.subscribers_[i].active.store(0, std::memory_order_release);
467 topic.subscribers_[i].owner_pid.store(0, std::memory_order_release);
468 topic.subscribers_[i].owner_starttime.store(0, std::memory_order_release);
469 topic.subscribers_[i].mode.store(
470 static_cast<uint32_t
>(LinuxSharedSubscriberMode::BROADCAST_FULL),
471 std::memory_order_release);
476 subscriber_index_ = i;
477 current_slot_index_ = INVALID_INDEX;
478 current_sequence_ = 0;
479 current_timestamp_ = MicrosecondTimestamp();
480 return ErrorCode::OK;
484 return ErrorCode::FULL;
487 LinuxSharedTopic* topic_ =
nullptr;
488 LinuxSharedTopic* owned_topic_ =
nullptr;
489 uint32_t subscriber_index_ = INVALID_INDEX;
490 uint32_t current_slot_index_ = INVALID_INDEX;
491 uint64_t current_sequence_ = 0;
492 MicrosecondTimestamp current_timestamp_;
510 SharedData() =
default;
516 ~SharedData() { Reset(); }
518 SharedData(
const SharedData&) =
delete;
519 SharedData& operator=(
const SharedData&) =
delete;
521 SharedData(SharedData&& other)
noexcept { *
this = std::move(other); }
526 SharedData& operator=(SharedData&& other)
noexcept
535 topic_ = other.topic_;
536 slot_index_ = other.slot_index_;
537 sequence_ = other.sequence_;
538 state_ = other.state_;
539 subscriber_index_ = other.subscriber_index_;
541 other.topic_ =
nullptr;
542 other.slot_index_ = INVALID_INDEX;
544 other.state_ = SharedDataState::EMPTY;
545 other.subscriber_index_ = INVALID_INDEX;
552 bool Valid()
const {
return topic_ !=
nullptr && slot_index_ != INVALID_INDEX; }
557 bool Empty()
const {
return !Valid(); }
562 uint64_t GetSequence()
const {
return sequence_; }
567 MicrosecondTimestamp GetTimestamp()
const
569 if (!Valid() || state_ != SharedDataState::SUBSCRIBER)
571 return MicrosecondTimestamp();
573 return topic_->SlotTimestamp(slot_index_);
587 return &topic_->payloads_[slot_index_];
595 TopicData* GetData()
const
601 return &topic_->payloads_[slot_index_];
614 if (state_ == SharedDataState::PUBLISHER)
616 topic_->RecycleSlot(slot_index_);
618 else if (state_ == SharedDataState::SUBSCRIBER)
620 topic_->ClearHeldSlot(subscriber_index_, slot_index_);
621 topic_->ReleaseSlot(slot_index_);
624 slot_index_ = INVALID_INDEX;
626 state_ = SharedDataState::EMPTY;
627 subscriber_index_ = INVALID_INDEX;
631 friend class LinuxSharedTopic<TopicData>;
632 friend class Subscriber;
634 LinuxSharedTopic* topic_ =
nullptr;
635 uint32_t slot_index_ = INVALID_INDEX;
636 uint64_t sequence_ = 0;
637 SharedDataState state_ = SharedDataState::EMPTY;
638 uint32_t subscriber_index_ = INVALID_INDEX;
645 explicit LinuxSharedTopic(
const char* topic_name)
646 : LinuxSharedTopic(topic_name, DEFAULT_DOMAIN_NAME)
650 LinuxSharedTopic(
const char* topic_name,
const char* domain_name)
654 domain_crc32_(ResolveDomainKey(domain_name)),
655 topic_name_(ResolveTopicName(topic_name)),
656 name_key_(BuildNameKey(domain_crc32_, topic_name_)),
657 shm_name_(BuildShmName(name_key_))
659 (void)ReadProcessIdentity(
static_cast<uint32_t
>(getpid()), self_identity_);
663 LinuxSharedTopic(
const char* topic_name, Topic::Domain& domain)
667 domain_crc32_(domain.node_ != nullptr ? domain.node_->key : 0),
668 topic_name_(ResolveTopicName(topic_name)),
669 name_key_(BuildNameKey(domain_crc32_, topic_name_)),
670 shm_name_(BuildShmName(name_key_))
672 (void)ReadProcessIdentity(
static_cast<uint32_t
>(getpid()), self_identity_);
682 LinuxSharedTopic(
const char* topic_name,
const LinuxSharedTopicConfig& config)
683 : LinuxSharedTopic(topic_name, DEFAULT_DOMAIN_NAME, config)
687 LinuxSharedTopic(
const char* topic_name,
const char* domain_name,
688 const LinuxSharedTopicConfig& config)
692 domain_crc32_(ResolveDomainKey(domain_name)),
693 topic_name_(ResolveTopicName(topic_name)),
694 name_key_(BuildNameKey(domain_crc32_, topic_name_)),
695 shm_name_(BuildShmName(name_key_))
697 (void)ReadProcessIdentity(
static_cast<uint32_t
>(getpid()), self_identity_);
701 LinuxSharedTopic(
const char* topic_name, Topic::Domain& domain,
702 const LinuxSharedTopicConfig& config)
706 domain_crc32_(domain.node_ != nullptr ? domain.node_->key : 0),
707 topic_name_(ResolveTopicName(topic_name)),
708 name_key_(BuildNameKey(domain_crc32_, topic_name_)),
709 shm_name_(BuildShmName(name_key_))
711 (void)ReadProcessIdentity(
static_cast<uint32_t
>(getpid()), self_identity_);
719 using SyncSubscriber = Subscriber;
724 ~LinuxSharedTopic() { Close(); }
726 LinuxSharedTopic(
const LinuxSharedTopic&) =
delete;
727 LinuxSharedTopic& operator=(
const LinuxSharedTopic&) =
delete;
729 LinuxSharedTopic(LinuxSharedTopic&&) =
delete;
730 LinuxSharedTopic& operator=(LinuxSharedTopic&&) =
delete;
736 bool Valid()
const {
return open_ok_; }
741 ErrorCode GetError()
const {
return open_status_; }
746 uint32_t GetSubscriberNum()
const
754 for (uint32_t i = 0; i < subscriber_capacity_; ++i)
756 if (subscribers_[i].active.load(std::memory_order_acquire) != 0)
774 return ErrorCode::STATE_ERR;
777 if (!PublisherValid())
779 return ErrorCode::STATE_ERR;
784 uint32_t slot_index = INVALID_INDEX;
785 ErrorCode pop_ans = PopFreeSlot(slot_index);
786 if (pop_ans != ErrorCode::OK)
788 ScavengeDeadSubscribers();
789 pop_ans = PopFreeSlot(slot_index);
790 if (pop_ans != ErrorCode::OK)
796 slots_[slot_index].refcount.store(0, std::memory_order_release);
797 slots_[slot_index].sequence.store(0, std::memory_order_release);
798 slots_[slot_index].timestamp_us = 0;
801 data.slot_index_ = slot_index;
803 data.state_ = SharedDataState::PUBLISHER;
804 data.subscriber_index_ = INVALID_INDEX;
805 return ErrorCode::OK;
815 SharedData topic_data;
816 const ErrorCode acquire_ans = CreateData(topic_data);
817 if (acquire_ans != ErrorCode::OK)
822 *topic_data.GetData() = data;
823 return Publish(topic_data);
826 ErrorCode Publish(
const TopicData& data, MicrosecondTimestamp timestamp)
828 SharedData topic_data;
829 const ErrorCode acquire_ans = CreateData(topic_data);
830 if (acquire_ans != ErrorCode::OK)
835 *topic_data.GetData() = data;
836 return Publish(topic_data, timestamp);
842 ErrorCode Publish(SharedData&& data) {
return PublishData<false>(data); }
844 ErrorCode Publish(SharedData&& data, MicrosecondTimestamp timestamp)
846 return PublishData<true>(data, timestamp);
852 ErrorCode Publish(SharedData& data) {
return PublishData<false>(data); }
854 ErrorCode Publish(SharedData& data, MicrosecondTimestamp timestamp)
856 return PublishData<true>(data, timestamp);
862 uint64_t GetPublishFailedNum()
const
868 return header_->publish_failures.load(std::memory_order_acquire);
876 static ErrorCode Remove(
const char* topic_name)
878 return Remove(topic_name, DEFAULT_DOMAIN_NAME);
881 static ErrorCode Remove(
const char* topic_name,
const char* domain_name)
883 const std::string shm_name = BuildShmName(
884 BuildNameKey(ResolveDomainKey(domain_name), ResolveTopicName(topic_name)));
885 if (shm_unlink(shm_name.c_str()) == 0 || errno == ENOENT)
887 return ErrorCode::OK;
889 return ErrorCode::FAILED;
892 static ErrorCode Remove(
const char* topic_name, Topic::Domain& domain)
894 const uint32_t domain_crc32 = (domain.node_ !=
nullptr) ? domain.node_->key : 0;
895 const std::string shm_name =
896 BuildShmName(BuildNameKey(domain_crc32, ResolveTopicName(topic_name)));
897 if (shm_unlink(shm_name.c_str()) == 0 || errno == ENOENT)
899 return ErrorCode::OK;
901 return ErrorCode::FAILED;
905 struct alignas(LibXR::CONCURRENCY_ALIGNMENT) SharedHeader
908 uint64_t name_key = 0;
909 uint32_t domain_crc32 = 0;
910 uint32_t version = 0;
911 uint32_t data_size = 0;
912 uint32_t slot_count = 0;
913 uint32_t subscriber_capacity = 0;
914 uint32_t queue_capacity = 0;
915 uint32_t topic_name_len = 0;
916 std::atomic<uint32_t> init_state;
917 std::atomic<uint32_t> publisher_pid;
918 std::atomic<uint64_t> publisher_starttime;
919 std::atomic<uint64_t> free_queue_head;
920 std::atomic<uint64_t> free_queue_tail;
921 std::atomic<uint64_t> next_sequence;
922 std::atomic<uint64_t> publish_failures;
925 struct alignas(LibXR::CONCURRENCY_ALIGNMENT) SlotControl
927 std::atomic<uint32_t> refcount;
928 std::atomic<uint64_t> sequence;
929 uint64_t timestamp_us;
932 struct alignas(16) FreeSlotCell
934 std::atomic<uint64_t> sequence;
935 uint32_t slot_index = 0;
936 uint32_t reserved = 0;
941 uint32_t slot_index = INVALID_INDEX;
942 uint32_t reserved = 0;
943 uint64_t sequence = 0;
946 struct alignas(LibXR::CONCURRENCY_ALIGNMENT) SubscriberControl
948 std::atomic<uint32_t> active;
949 std::atomic<uint32_t> mode;
950 std::atomic<uint32_t> queue_head;
951 std::atomic<uint32_t> queue_tail;
952 std::atomic<uint32_t> ready_sem_count;
953 std::atomic<uint64_t> dropped_messages;
954 std::atomic<uint32_t> owner_pid;
955 std::atomic<uint64_t> owner_starttime;
956 std::atomic<uint32_t> held_slot;
959 struct alignas(LibXR::CONCURRENCY_ALIGNMENT) BalancedGroupControl
961 std::atomic<uint64_t> rr_cursor;
964 struct ProcessIdentity
967 uint64_t starttime = 0;
970 static constexpr uint64_t MAGIC = 0x4c58524950435348ULL;
971 static constexpr uint32_t VERSION = 2;
972 static constexpr uint32_t INIT_READY = 1;
973 static constexpr uint32_t INVALID_INDEX = UINT32_MAX;
975 static uint32_t ResolveDomainKey(
const char* domain_name)
977 const std::string resolved = (domain_name ==
nullptr || domain_name[0] ==
'\0')
978 ? std::string(DEFAULT_DOMAIN_NAME)
979 : std::string(domain_name);
980 return CRC32::Calculate(resolved.data(), resolved.size());
983 static std::string ResolveTopicName(
const char* topic_name)
985 return (topic_name !=
nullptr) ? std::string(topic_name) : std::string();
988 static uint64_t BuildNameKey(uint32_t domain_crc32,
const std::string& topic_name)
990 const uint32_t topic_len =
static_cast<uint32_t
>(topic_name.size());
991 std::string key_material;
992 key_material.reserve(
sizeof(domain_crc32) +
sizeof(topic_len) + topic_len);
993 key_material.append(
reinterpret_cast<const char*
>(&domain_crc32),
994 sizeof(domain_crc32));
995 key_material.append(
reinterpret_cast<const char*
>(&topic_len),
sizeof(topic_len));
996 key_material.append(topic_name.data(), topic_name.size());
997 return CRC64::Calculate(key_material.data(), key_material.size());
1000 static std::string BuildShmName(uint64_t name_key)
1002 char buffer[64] = {};
1003 std::snprintf(buffer,
sizeof(buffer),
"/libxr_ipc_%016" PRIx64, name_key);
1004 return std::string(buffer);
1007 static size_t AlignUp(
size_t value,
size_t alignment)
1009 return (value + alignment - 1U) & ~(alignment - 1U);
1012 static uint64_t NowMonotonicMs() {
return MonotonicTime::NowMilliseconds(); }
1014 static MicrosecondTimestamp NowMessageTimestamp() {
return Topic::NowTimestamp(); }
1016 static uint64_t ToSharedTimestamp(MicrosecondTimestamp timestamp)
1018 return MonotonicTime::XrToSharedMicroseconds(
static_cast<uint64_t
>(timestamp));
1021 static MicrosecondTimestamp FromSharedTimestamp(uint64_t timestamp_us)
1023 return MicrosecondTimestamp(MonotonicTime::SharedToXrMicroseconds(timestamp_us));
1026 static bool ReadProcessIdentity(uint32_t pid, ProcessIdentity& identity)
1035 std::snprintf(path,
sizeof(path),
"/proc/%u/stat", pid);
1037 std::ifstream file(path);
1038 if (!file.is_open())
1044 std::getline(file, line);
1050 const size_t rparen = line.rfind(
')');
1051 if (rparen == std::string::npos || rparen + 2U >= line.size())
1056 std::istringstream iss(line.substr(rparen + 2U));
1058 for (
int field = 3; field <= 22; ++field)
1060 if (!(iss >> token))
1068 identity.starttime = std::strtoull(token.c_str(),
nullptr, 10);
1069 return identity.starttime != 0;
1076 static int FutexWait(std::atomic<uint32_t>* word, uint32_t expected,
1077 uint32_t timeout_ms)
1079 struct timespec timeout = {};
1080 struct timespec* timeout_ptr =
nullptr;
1081 if (timeout_ms != UINT32_MAX)
1083 timeout.tv_sec =
static_cast<time_t
>(timeout_ms / 1000U);
1084 timeout.tv_nsec =
static_cast<long>(timeout_ms % 1000U) * 1000000L;
1085 timeout_ptr = &timeout;
1088 return static_cast<int>(syscall(SYS_futex,
reinterpret_cast<uint32_t*
>(word),
1089 FUTEX_WAIT, expected, timeout_ptr,
nullptr, 0));
1092 static int FutexWake(std::atomic<uint32_t>* word)
1094 return static_cast<int>(syscall(SYS_futex,
reinterpret_cast<uint32_t*
>(word),
1095 FUTEX_WAKE, INT32_MAX,
nullptr,
nullptr, 0));
1098 static size_t ComputeSharedBytes(uint32_t slot_count, uint32_t subscriber_capacity,
1099 uint32_t queue_capacity, uint32_t topic_name_len)
1102 offset = AlignUp(offset,
alignof(SharedHeader));
1103 offset +=
sizeof(SharedHeader);
1105 offset +=
static_cast<size_t>(topic_name_len) + 1U;
1107 offset = AlignUp(offset,
alignof(SlotControl));
1108 offset +=
sizeof(SlotControl) * slot_count;
1110 offset = AlignUp(offset,
alignof(SubscriberControl));
1111 offset +=
sizeof(SubscriberControl) * subscriber_capacity;
1113 offset = AlignUp(offset,
alignof(BalancedGroupControl));
1114 offset +=
sizeof(BalancedGroupControl);
1116 offset = AlignUp(offset,
alignof(std::atomic<uint32_t>));
1117 offset +=
sizeof(std::atomic<uint32_t>) * subscriber_capacity;
1119 offset = AlignUp(offset,
alignof(FreeSlotCell));
1120 offset +=
sizeof(FreeSlotCell) * slot_count;
1122 offset = AlignUp(offset,
alignof(Descriptor));
1123 offset +=
sizeof(Descriptor) * subscriber_capacity * queue_capacity;
1125 offset = AlignUp(offset,
alignof(TopicData));
1126 offset +=
sizeof(TopicData) * slot_count;
1130 void SetupPointers()
1134 offset = AlignUp(offset,
alignof(SharedHeader));
1135 header_ =
reinterpret_cast<SharedHeader*
>(base_ + offset);
1136 offset +=
sizeof(SharedHeader);
1138 topic_name_ptr_ =
reinterpret_cast<char*
>(base_ + offset);
1139 offset +=
static_cast<size_t>(header_->topic_name_len) + 1U;
1141 offset = AlignUp(offset,
alignof(SlotControl));
1142 slots_ =
reinterpret_cast<SlotControl*
>(base_ + offset);
1143 offset +=
sizeof(SlotControl) * slot_count_;
1145 offset = AlignUp(offset,
alignof(SubscriberControl));
1146 subscribers_ =
reinterpret_cast<SubscriberControl*
>(base_ + offset);
1147 offset +=
sizeof(SubscriberControl) * subscriber_capacity_;
1149 offset = AlignUp(offset,
alignof(BalancedGroupControl));
1150 balanced_group_ =
reinterpret_cast<BalancedGroupControl*
>(base_ + offset);
1151 offset +=
sizeof(BalancedGroupControl);
1153 offset = AlignUp(offset,
alignof(std::atomic<uint32_t>));
1154 balanced_members_ =
reinterpret_cast<std::atomic<uint32_t>*
>(base_ + offset);
1155 offset +=
sizeof(std::atomic<uint32_t>) * subscriber_capacity_;
1157 offset = AlignUp(offset,
alignof(FreeSlotCell));
1158 free_slots_ =
reinterpret_cast<FreeSlotCell*
>(base_ + offset);
1159 offset +=
sizeof(FreeSlotCell) * slot_count_;
1161 offset = AlignUp(offset,
alignof(Descriptor));
1162 descriptors_ =
reinterpret_cast<Descriptor*
>(base_ + offset);
1163 offset +=
sizeof(Descriptor) * subscriber_capacity_ * queue_capacity_;
1165 offset = AlignUp(offset,
alignof(TopicData));
1166 payloads_ =
reinterpret_cast<TopicData*
>(base_ + offset);
1169 bool HeaderMatchesIdentity()
const
1171 if (header_->name_key != name_key_)
1175 if (header_->domain_crc32 != domain_crc32_)
1179 if (header_->topic_name_len != topic_name_.size())
1183 if (std::memcmp(topic_name_ptr_, topic_name_.c_str(), topic_name_.size() + 1U) != 0)
1192 if (config_.slot_num == 0 || config_.subscriber_num == 0 || config_.queue_num < 2)
1194 return ErrorCode::ARG_ERR;
1197 const size_t bytes =
1198 ComputeSharedBytes(config_.slot_num, config_.subscriber_num, config_.queue_num,
1199 static_cast<uint32_t
>(topic_name_.size()));
1201 if (ftruncate(fd_,
static_cast<off_t
>(bytes)) != 0)
1203 return ErrorCode::INIT_ERR;
1206 const struct stat st = GetStat();
1207 if (st.st_size <= 0)
1209 return ErrorCode::INIT_ERR;
1212 mapping_size_ =
static_cast<size_t>(st.st_size);
1213 mapping_ = mmap(
nullptr, mapping_size_, PROT_READ | PROT_WRITE, MAP_SHARED, fd_, 0);
1214 if (mapping_ == MAP_FAILED)
1217 return ErrorCode::INIT_ERR;
1220 base_ =
static_cast<uint8_t*
>(mapping_);
1221 slot_count_ = config_.slot_num;
1222 subscriber_capacity_ = config_.subscriber_num;
1223 queue_capacity_ = config_.queue_num;
1224 header_ =
reinterpret_cast<SharedHeader*
>(base_ + AlignUp(0,
alignof(SharedHeader)));
1225 header_->topic_name_len =
static_cast<uint32_t
>(topic_name_.size());
1228 header_->magic = MAGIC;
1229 header_->name_key = name_key_;
1230 header_->domain_crc32 = domain_crc32_;
1231 header_->version = VERSION;
1232 header_->data_size =
sizeof(TopicData);
1233 header_->slot_count = slot_count_;
1234 header_->subscriber_capacity = subscriber_capacity_;
1235 header_->queue_capacity = queue_capacity_;
1236 std::memcpy(topic_name_ptr_, topic_name_.c_str(), topic_name_.size() + 1U);
1237 header_->publisher_pid.store(self_identity_.pid, std::memory_order_release);
1238 header_->publisher_starttime.store(self_identity_.starttime,
1239 std::memory_order_release);
1240 header_->free_queue_head.store(0, std::memory_order_release);
1241 header_->free_queue_tail.store(slot_count_, std::memory_order_release);
1242 header_->next_sequence.store(0, std::memory_order_release);
1243 header_->publish_failures.store(0, std::memory_order_release);
1245 for (uint32_t i = 0; i < slot_count_; ++i)
1247 slots_[i].refcount.store(0, std::memory_order_release);
1248 slots_[i].sequence.store(0, std::memory_order_release);
1249 slots_[i].timestamp_us = 0;
1250 std::construct_at(&payloads_[i], TopicData{});
1251 free_slots_[i].slot_index = i;
1252 free_slots_[i].sequence.store(
static_cast<uint64_t
>(i) + 1U,
1253 std::memory_order_release);
1256 for (uint32_t i = 0; i < subscriber_capacity_; ++i)
1258 subscribers_[i].active.store(0, std::memory_order_release);
1259 subscribers_[i].mode.store(
1260 static_cast<uint32_t
>(LinuxSharedSubscriberMode::BROADCAST_FULL),
1261 std::memory_order_release);
1262 subscribers_[i].queue_head.store(0, std::memory_order_release);
1263 subscribers_[i].queue_tail.store(0, std::memory_order_release);
1264 subscribers_[i].ready_sem_count.store(0, std::memory_order_release);
1265 subscribers_[i].dropped_messages.store(0, std::memory_order_release);
1266 subscribers_[i].owner_pid.store(0, std::memory_order_release);
1267 subscribers_[i].owner_starttime.store(0, std::memory_order_release);
1268 subscribers_[i].held_slot.store(INVALID_INDEX, std::memory_order_release);
1269 balanced_members_[i].store(INVALID_INDEX, std::memory_order_release);
1272 balanced_group_->rr_cursor.store(0, std::memory_order_release);
1274 for (
size_t i = 0; i < static_cast<size_t>(subscriber_capacity_) * queue_capacity_;
1277 descriptors_[i] = Descriptor{};
1280 header_->init_state.store(INIT_READY, std::memory_order_release);
1281 return ErrorCode::OK;
1286 const struct stat st = GetStat();
1287 if (st.st_size <= 0)
1289 return ErrorCode::NOT_FOUND;
1292 mapping_size_ =
static_cast<size_t>(st.st_size);
1293 mapping_ = mmap(
nullptr, mapping_size_, PROT_READ | PROT_WRITE, MAP_SHARED, fd_, 0);
1294 if (mapping_ == MAP_FAILED)
1297 return ErrorCode::INIT_ERR;
1300 base_ =
static_cast<uint8_t*
>(mapping_);
1301 header_ =
reinterpret_cast<SharedHeader*
>(base_);
1303 while (header_->init_state.load(std::memory_order_acquire) != INIT_READY)
1308 if (header_->magic != MAGIC || header_->version != VERSION ||
1309 header_->data_size !=
sizeof(TopicData))
1311 return ErrorCode::CHECK_ERR;
1314 slot_count_ = header_->slot_count;
1315 subscriber_capacity_ = header_->subscriber_capacity;
1316 queue_capacity_ = header_->queue_capacity;
1318 if (!HeaderMatchesIdentity())
1320 return ErrorCode::CHECK_ERR;
1322 return ErrorCode::OK;
1325 bool TryReclaimStaleSegment()
1327 int stale_fd = shm_open(shm_name_.c_str(), O_RDWR, 0600);
1330 return errno == ENOENT;
1333 struct stat st = {};
1334 if (fstat(stale_fd, &st) != 0)
1340 bool reclaim =
false;
1341 if (st.st_size <
static_cast<off_t
>(
sizeof(SharedHeader)))
1347 void* mapping = mmap(
nullptr,
static_cast<size_t>(st.st_size),
1348 PROT_READ | PROT_WRITE, MAP_SHARED, stale_fd, 0);
1349 if (mapping != MAP_FAILED)
1351 uint8_t* base =
static_cast<uint8_t*
>(mapping);
1352 auto* header =
reinterpret_cast<SharedHeader*
>(mapping);
1353 const uint32_t init_state = header->init_state.load(std::memory_order_acquire);
1354 bool identity_match =
false;
1355 const size_t mapping_size =
static_cast<size_t>(st.st_size);
1356 const size_t topic_name_bytes = topic_name_.size() + 1U;
1357 if (header->magic == MAGIC && header->version == VERSION &&
1358 header->domain_crc32 == domain_crc32_ &&
1359 header->topic_name_len == topic_name_.size())
1361 size_t offset = AlignUp(0,
alignof(SharedHeader));
1362 offset +=
sizeof(SharedHeader);
1363 if (offset <= mapping_size && topic_name_bytes <= (mapping_size - offset))
1365 const char* mapped_topic =
reinterpret_cast<const char*
>(base + offset);
1367 (header->name_key == name_key_) &&
1368 (std::memcmp(mapped_topic, topic_name_.c_str(), topic_name_bytes) == 0);
1371 const ProcessIdentity publisher_identity = {
1372 header->publisher_pid.load(std::memory_order_acquire),
1373 header->publisher_starttime.load(std::memory_order_acquire),
1376 if (!identity_match)
1380 else if (init_state != INIT_READY)
1382 reclaim = !ProcessAlive(publisher_identity);
1384 else if (!ProcessAlive(publisher_identity))
1389 munmap(mapping,
static_cast<size_t>(st.st_size));
1400 return shm_unlink(shm_name_.c_str()) == 0 || errno == ENOENT;
1405 open_status_ = ErrorCode::STATE_ERR;
1410 for (
int attempt = 0; attempt < 2; ++attempt)
1412 fd_ = shm_open(shm_name_.c_str(), O_CREAT | O_EXCL | O_RDWR, 0600);
1418 if (errno != EEXIST || !TryReclaimStaleSegment())
1430 open_status_ = InitializeLayout();
1434 fd_ = shm_open(shm_name_.c_str(), O_RDWR, 0600);
1441 open_status_ = AttachLayout();
1450 open_ok_ = (open_status_ == ErrorCode::OK);
1455 if (mapping_ !=
nullptr)
1457 munmap(mapping_, mapping_size_);
1464 subscribers_ =
nullptr;
1465 free_slots_ =
nullptr;
1466 descriptors_ =
nullptr;
1467 payloads_ =
nullptr;
1470 open_status_ = ErrorCode::STATE_ERR;
1473 struct stat GetStat() const
1475 struct stat st = {};
1480 Descriptor* DescriptorRing(uint32_t subscriber_index)
const
1482 return descriptors_ +
static_cast<size_t>(subscriber_index) * queue_capacity_;
1485 static bool ProcessAlive(
const ProcessIdentity& identity)
1487 ProcessIdentity current = {};
1488 if (!ReadProcessIdentity(identity.pid, current))
1493 return current.starttime == identity.starttime;
1496 bool PublisherValid()
const
1498 if (!publisher_ || header_ ==
nullptr)
1503 const ProcessIdentity owner = {
1504 header_->publisher_pid.load(std::memory_order_acquire),
1505 header_->publisher_starttime.load(std::memory_order_acquire),
1507 return owner.pid == self_identity_.pid && owner.starttime == self_identity_.starttime;
1510 void HoldSlot(uint32_t subscriber_index, uint32_t slot_index)
1512 subscribers_[subscriber_index].held_slot.store(slot_index, std::memory_order_release);
1515 void ClearHeldSlot(uint32_t subscriber_index, uint32_t slot_index)
1517 uint32_t expected = slot_index;
1518 subscribers_[subscriber_index].held_slot.compare_exchange_strong(
1519 expected, INVALID_INDEX, std::memory_order_acq_rel, std::memory_order_relaxed);
1522 MicrosecondTimestamp SlotTimestamp(uint32_t slot_index)
const
1524 return FromSharedTimestamp(slots_[slot_index].timestamp_us);
1527 ErrorCode RegisterBalancedSubscriber(uint32_t subscriber_index)
1529 for (uint32_t i = 0; i < subscriber_capacity_; ++i)
1531 uint32_t expected = INVALID_INDEX;
1532 if (balanced_members_[i].compare_exchange_strong(expected, subscriber_index,
1533 std::memory_order_acq_rel,
1534 std::memory_order_relaxed))
1536 return ErrorCode::OK;
1539 return ErrorCode::FULL;
1542 void UnregisterBalancedSubscriber(uint32_t subscriber_index)
1544 for (uint32_t i = 0; i < subscriber_capacity_; ++i)
1546 uint32_t expected = subscriber_index;
1547 if (balanced_members_[i].compare_exchange_strong(expected, INVALID_INDEX,
1548 std::memory_order_acq_rel,
1549 std::memory_order_relaxed))
1556 bool SelectBalancedSubscriber(uint32_t& subscriber_index)
1558 const uint64_t base =
1559 balanced_group_->rr_cursor.fetch_add(1, std::memory_order_acq_rel);
1560 for (uint32_t offset = 0; offset < subscriber_capacity_; ++offset)
1562 const uint32_t member_index =
1563 balanced_members_[(base + offset) % subscriber_capacity_].load(
1564 std::memory_order_acquire);
1565 if (member_index == INVALID_INDEX)
1569 if (subscribers_[member_index].active.load(std::memory_order_acquire) == 0)
1573 if (subscribers_[member_index].mode.load(std::memory_order_acquire) !=
1574 static_cast<uint32_t
>(LinuxSharedSubscriberMode::BALANCE_RR))
1578 const ProcessIdentity owner_identity = {
1579 subscribers_[member_index].owner_pid.load(std::memory_order_acquire),
1580 subscribers_[member_index].owner_starttime.load(std::memory_order_acquire),
1582 if (owner_identity.pid == 0 || owner_identity.starttime == 0)
1586 if (!ProcessAlive(owner_identity))
1588 ReclaimSubscriber(member_index);
1591 if (!QueueHasSpace(member_index))
1595 subscriber_index = member_index;
1601 bool ReclaimSubscriber(uint32_t subscriber_index)
1603 uint32_t expected = 1;
1604 if (!subscribers_[subscriber_index].active.compare_exchange_strong(
1605 expected, 0, std::memory_order_acq_rel, std::memory_order_relaxed))
1610 if (subscribers_[subscriber_index].mode.load(std::memory_order_acquire) ==
1611 static_cast<uint32_t
>(LinuxSharedSubscriberMode::BALANCE_RR))
1613 UnregisterBalancedSubscriber(subscriber_index);
1616 subscribers_[subscriber_index].owner_pid.store(0, std::memory_order_release);
1617 subscribers_[subscriber_index].owner_starttime.store(0, std::memory_order_release);
1618 subscribers_[subscriber_index].mode.store(
1619 static_cast<uint32_t
>(LinuxSharedSubscriberMode::BROADCAST_FULL),
1620 std::memory_order_release);
1622 const uint32_t held_slot = subscribers_[subscriber_index].held_slot.exchange(
1623 INVALID_INDEX, std::memory_order_acq_rel);
1624 if (held_slot != INVALID_INDEX)
1626 ReleaseSlot(held_slot);
1629 Descriptor desc = {};
1630 while (TryPopDescriptor(subscriber_index, desc) == ErrorCode::OK)
1632 ReleaseSlot(desc.slot_index);
1638 static void PostReady(SubscriberControl& control)
1640 control.ready_sem_count.fetch_add(1, std::memory_order_release);
1641 FutexWake(&control.ready_sem_count);
1644 static void ConsumeReady(SubscriberControl& control)
1646 const uint32_t prev = control.ready_sem_count.fetch_sub(1, std::memory_order_acq_rel);
1650 static ErrorCode WaitReady(SubscriberControl& control, uint32_t timeout_ms)
1652 if (control.ready_sem_count.load(std::memory_order_acquire) != 0)
1654 return ErrorCode::OK;
1657 const bool infinite_wait = (timeout_ms == UINT32_MAX);
1658 const uint64_t deadline_ms = infinite_wait ? 0 : (NowMonotonicMs() + timeout_ms);
1662 if (control.ready_sem_count.load(std::memory_order_acquire) != 0)
1664 return ErrorCode::OK;
1667 uint32_t wait_ms = UINT32_MAX;
1670 wait_ms = MonotonicTime::RemainingMilliseconds(deadline_ms);
1673 return ErrorCode::TIMEOUT;
1677 wait_ms = MonotonicTime::WaitSliceMilliseconds(wait_ms);
1679 const int futex_ans = FutexWait(&control.ready_sem_count, 0, wait_ms);
1680 if (futex_ans == 0 || errno == EAGAIN || errno == EINTR)
1685 if (errno == ETIMEDOUT)
1691 if (MonotonicTime::RemainingMilliseconds(deadline_ms) == 0 &&
1692 control.ready_sem_count.load(std::memory_order_acquire) == 0)
1694 return ErrorCode::TIMEOUT;
1699 return ErrorCode::FAILED;
1703 void ScavengeDeadSubscribers()
1705 for (uint32_t i = 0; i < subscriber_capacity_; ++i)
1707 if (subscribers_[i].active.load(std::memory_order_acquire) == 0)
1712 const ProcessIdentity owner_identity = {
1713 subscribers_[i].owner_pid.load(std::memory_order_acquire),
1714 subscribers_[i].owner_starttime.load(std::memory_order_acquire),
1716 if (ProcessAlive(owner_identity))
1721 ReclaimSubscriber(i);
1725 bool QueueHasSpace(uint32_t subscriber_index)
const
1727 const SubscriberControl& control = subscribers_[subscriber_index];
1728 const uint32_t head = control.queue_head.load(std::memory_order_acquire);
1729 const uint32_t tail = control.queue_tail.load(std::memory_order_relaxed);
1730 const uint32_t next_tail = (tail + 1U) % queue_capacity_;
1731 return next_tail != head;
1734 void PushDescriptor(uint32_t subscriber_index,
const Descriptor& descriptor)
1736 SubscriberControl& control = subscribers_[subscriber_index];
1737 Descriptor* ring = DescriptorRing(subscriber_index);
1739 const uint32_t tail = control.queue_tail.load(std::memory_order_relaxed);
1740 ring[tail] = descriptor;
1741 const uint32_t next_tail = (tail + 1U) % queue_capacity_;
1742 control.queue_tail.store(next_tail, std::memory_order_release);
1746 ErrorCode TryPopDescriptor(uint32_t subscriber_index, Descriptor& descriptor)
1748 SubscriberControl& control = subscribers_[subscriber_index];
1749 Descriptor* ring = DescriptorRing(subscriber_index);
1753 uint32_t head = control.queue_head.load(std::memory_order_relaxed);
1754 const uint32_t tail = control.queue_tail.load(std::memory_order_acquire);
1757 return ErrorCode::EMPTY;
1760 descriptor = ring[head];
1761 const uint32_t next_head = (head + 1U) % queue_capacity_;
1762 if (control.queue_head.compare_exchange_weak(
1763 head, next_head, std::memory_order_acq_rel, std::memory_order_relaxed))
1765 ConsumeReady(control);
1766 return ErrorCode::OK;
1771 ErrorCode DropDescriptor(uint32_t subscriber_index)
1773 SubscriberControl& control = subscribers_[subscriber_index];
1774 Descriptor* ring = DescriptorRing(subscriber_index);
1778 uint32_t head = control.queue_head.load(std::memory_order_relaxed);
1779 const uint32_t tail = control.queue_tail.load(std::memory_order_acquire);
1782 return ErrorCode::EMPTY;
1785 const Descriptor descriptor = ring[head];
1786 const uint32_t next_head = (head + 1U) % queue_capacity_;
1787 if (control.queue_head.compare_exchange_weak(
1788 head, next_head, std::memory_order_acq_rel, std::memory_order_relaxed))
1790 control.dropped_messages.fetch_add(1, std::memory_order_relaxed);
1791 ConsumeReady(control);
1792 ReleaseSlot(descriptor.slot_index);
1793 return ErrorCode::OK;
1798 ErrorCode PopFreeSlot(uint32_t& slot_index)
1802 uint64_t head = header_->free_queue_head.load(std::memory_order_relaxed);
1803 FreeSlotCell& cell = free_slots_[head % slot_count_];
1804 const uint64_t seq = cell.sequence.load(std::memory_order_acquire);
1805 const intptr_t diff =
static_cast<intptr_t
>(seq) -
static_cast<intptr_t
>(head + 1U);
1809 if (header_->free_queue_head.compare_exchange_weak(
1810 head, head + 1U, std::memory_order_acq_rel, std::memory_order_relaxed))
1812 slot_index = cell.slot_index;
1813 cell.sequence.store(head + slot_count_, std::memory_order_release);
1814 return ErrorCode::OK;
1819 return ErrorCode::FULL;
1824 void RecycleSlot(uint32_t slot_index)
1826 slots_[slot_index].sequence.store(0, std::memory_order_release);
1827 slots_[slot_index].timestamp_us = 0;
1831 uint64_t tail = header_->free_queue_tail.load(std::memory_order_relaxed);
1832 FreeSlotCell& cell = free_slots_[tail % slot_count_];
1833 const uint64_t seq = cell.sequence.load(std::memory_order_acquire);
1834 const intptr_t diff =
static_cast<intptr_t
>(seq) -
static_cast<intptr_t
>(tail);
1838 if (header_->free_queue_tail.compare_exchange_weak(
1839 tail, tail + 1U, std::memory_order_acq_rel, std::memory_order_relaxed))
1841 cell.slot_index = slot_index;
1842 cell.sequence.store(tail + 1U, std::memory_order_release);
1849 void ReleaseSlot(uint32_t slot_index)
1851 const uint32_t prev =
1852 slots_[slot_index].refcount.fetch_sub(1, std::memory_order_acq_rel);
1856 RecycleSlot(slot_index);
1860 template <
bool HAS_TIMESTAMP>
1862 MicrosecondTimestamp timestamp = MicrosecondTimestamp())
1864 if (!data.Valid() || data.topic_ !=
this)
1866 return ErrorCode::STATE_ERR;
1869 uint32_t active_count = 0;
1870 uint32_t balanced_target = INVALID_INDEX;
1871 bool has_balanced_subscriber =
false;
1872 for (uint32_t i = 0; i < subscriber_capacity_; ++i)
1874 if (subscribers_[i].active.load(std::memory_order_acquire) == 0)
1879 const LinuxSharedSubscriberMode mode =
static_cast<LinuxSharedSubscriberMode
>(
1880 subscribers_[i].mode.load(std::memory_order_acquire));
1881 if (mode == LinuxSharedSubscriberMode::BALANCE_RR)
1883 has_balanced_subscriber =
true;
1887 if (!QueueHasSpace(i))
1889 ScavengeDeadSubscribers();
1890 if (subscribers_[i].active.load(std::memory_order_acquire) == 0)
1895 if (mode == LinuxSharedSubscriberMode::BROADCAST_DROP_OLD)
1897 const ErrorCode drop_ans = DropDescriptor(i);
1898 if (drop_ans == ErrorCode::EMPTY && QueueHasSpace(i))
1903 else if (drop_ans != ErrorCode::OK)
1905 header_->publish_failures.fetch_add(1, std::memory_order_relaxed);
1907 return ErrorCode::FULL;
1912 subscribers_[i].dropped_messages.fetch_add(1, std::memory_order_relaxed);
1913 header_->publish_failures.fetch_add(1, std::memory_order_relaxed);
1915 return ErrorCode::FULL;
1922 if (has_balanced_subscriber)
1924 if (!SelectBalancedSubscriber(balanced_target))
1926 ScavengeDeadSubscribers();
1927 if (!SelectBalancedSubscriber(balanced_target))
1929 header_->publish_failures.fetch_add(1, std::memory_order_relaxed);
1931 return ErrorCode::FULL;
1937 if (active_count == 0)
1940 return ErrorCode::OK;
1943 const uint64_t sequence =
1944 header_->next_sequence.fetch_add(1, std::memory_order_acq_rel) + 1ULL;
1945 if constexpr (!HAS_TIMESTAMP)
1947 timestamp = NowMessageTimestamp();
1949 SlotControl& slot = slots_[data.slot_index_];
1950 slot.refcount.store(active_count, std::memory_order_release);
1951 slot.timestamp_us = ToSharedTimestamp(timestamp);
1952 slot.sequence.store(sequence, std::memory_order_release);
1954 const Descriptor descriptor = {data.slot_index_, 0U, sequence};
1955 for (uint32_t i = 0; i < subscriber_capacity_; ++i)
1957 if (subscribers_[i].active.load(std::memory_order_acquire) == 0)
1961 const LinuxSharedSubscriberMode mode =
static_cast<LinuxSharedSubscriberMode
>(
1962 subscribers_[i].mode.load(std::memory_order_acquire));
1963 if (mode == LinuxSharedSubscriberMode::BALANCE_RR)
1967 PushDescriptor(i, descriptor);
1970 if (balanced_target != INVALID_INDEX)
1972 PushDescriptor(balanced_target, descriptor);
1975 data.topic_ =
nullptr;
1976 data.slot_index_ = INVALID_INDEX;
1977 return ErrorCode::OK;
1980 bool create_ =
false;
1981 bool publisher_ =
false;
1982 LinuxSharedTopicConfig config_;
1983 uint32_t domain_crc32_ = 0;
1984 std::string topic_name_;
1985 uint64_t name_key_ = 0;
1986 std::string shm_name_;
1987 ProcessIdentity self_identity_ = {};
1990 void* mapping_ =
nullptr;
1991 uint8_t* base_ =
nullptr;
1992 size_t mapping_size_ = 0;
1994 SharedHeader* header_ =
nullptr;
1995 char* topic_name_ptr_ =
nullptr;
1996 SlotControl* slots_ =
nullptr;
1997 SubscriberControl* subscribers_ =
nullptr;
1998 BalancedGroupControl* balanced_group_ =
nullptr;
1999 std::atomic<uint32_t>* balanced_members_ =
nullptr;
2000 FreeSlotCell* free_slots_ =
nullptr;
2001 Descriptor* descriptors_ =
nullptr;
2002 TopicData* payloads_ =
nullptr;
2004 uint32_t slot_count_ = 0;
2005 uint32_t subscriber_capacity_ = 0;
2006 uint32_t queue_capacity_ = 0;
2008 bool open_ok_ =
false;
2009 ErrorCode open_status_ = ErrorCode::STATE_ERR;
@ INIT_ERR
初始化错误 | Initialization error