6 const std::shared_ptr<DataProcessor>& processor,
7 const std::shared_ptr<SignalSourceContainer>& container,
8 const std::string& tag)
14 processor->on_attach(container);
18 const std::shared_ptr<DataProcessor>& processor,
19 const std::shared_ptr<SignalSourceContainer>& container,
24 if (position >= processors.size()) {
25 processors.push_back(processor);
27 processors.insert(processors.begin() +
static_cast<std::ptrdiff_t
>(position), processor);
30 processor->on_attach(container);
34 const std::shared_ptr<DataProcessor>& processor,
35 const std::shared_ptr<SignalSourceContainer>& container)
46 const std::shared_ptr<DataProcessor>& processor,
47 const std::shared_ptr<SignalSourceContainer>& container)
54 auto& processors = it->second;
55 auto proc_it = std::ranges::find(processors, processor);
57 if (proc_it != processors.end()) {
58 processor->on_detach(container);
59 processors.erase(proc_it);
62 if (processors.empty()) {
78 bool expected =
false;
80 std::memory_order_acquire, std::memory_order_relaxed)) {
86 for (
auto& processor : it->second) {
87 processor->process(container);
96 const std::shared_ptr<SignalSourceContainer>& container,
97 const std::function<
bool(
const std::shared_ptr<DataProcessor>&)>& filter)
99 bool expected =
false;
101 std::memory_order_acquire, std::memory_order_relaxed)) {
107 for (
auto& processor : it->second) {
108 if (filter(processor)) {
109 processor->process(container);
119 const std::shared_ptr<SignalSourceContainer>& container,
120 const std::string& tag)
122 bool expected =
false;
124 std::memory_order_acquire, std::memory_order_relaxed)) {
130 for (
auto& processor : it->second) {
133 processor->process(container);
void add_processor_at(const std::shared_ptr< DataProcessor > &processor, const std::shared_ptr< SignalSourceContainer > &container, size_t position)
Adds a processor at a specific position in the chain.
void process_tagged(const std::shared_ptr< SignalSourceContainer > &container, const std::string &tag)
Processes a container with only the processors carrying a specific tag.
void add_processor(const std::shared_ptr< DataProcessor > &processor, const std::shared_ptr< SignalSourceContainer > &container, const std::string &tag="")
Adds a processor to the end of the chain for a specific container.
void process(const std::shared_ptr< SignalSourceContainer > &container)
Processes a container with all its associated processors in sequence.
void remove_processor(const std::shared_ptr< DataProcessor > &processor, const std::shared_ptr< SignalSourceContainer > &container)
Removes a processor from a container's chain.
void drain_pending_removals()
Flushes all removals deferred during the most recent process iteration.
void remove_processor_direct(const std::shared_ptr< DataProcessor > &processor, const std::shared_ptr< SignalSourceContainer > &container)
Performs immediate removal of a processor from the chain.
void process_filtered(const std::shared_ptr< SignalSourceContainer > &container, const std::function< bool(const std::shared_ptr< DataProcessor > &)> &filter)
Processes a container with only the processors matching a filter predicate.
std::unordered_map< std::shared_ptr< SignalSourceContainer >, std::vector< std::shared_ptr< DataProcessor > > > m_container_processors
Maps each container to its ordered processor sequence.
std::vector< std::pair< std::shared_ptr< DataProcessor >, std::shared_ptr< SignalSourceContainer > > > m_pending_removal
Removals deferred because they arrived during active iteration.
std::atomic< bool > m_is_processing
Guards all process variants against concurrent or re-entrant iteration.
std::unordered_map< std::shared_ptr< DataProcessor >, std::string > m_processor_tags
Maps processors to their optional tag strings.