MayaFlux 0.5.0
Digital-First Multimedia Processing Framework
Loading...
Searching...
No Matches
DataProcessingChain.cpp
Go to the documentation of this file.
2
3namespace MayaFlux::Kakshya {
4
6 const std::shared_ptr<DataProcessor>& processor,
7 const std::shared_ptr<SignalSourceContainer>& container,
8 const std::string& tag)
9{
10 m_container_processors[container].emplace_back(processor);
11 if (!tag.empty()) {
12 m_processor_tags[processor] = tag;
13 }
14 processor->on_attach(container);
15}
16
18 const std::shared_ptr<DataProcessor>& processor,
19 const std::shared_ptr<SignalSourceContainer>& container,
20 size_t position)
21{
22 auto& processors = m_container_processors[container];
23
24 if (position >= processors.size()) {
25 processors.push_back(processor);
26 } else {
27 processors.insert(processors.begin() + static_cast<std::ptrdiff_t>(position), processor);
28 }
29
30 processor->on_attach(container);
31}
32
34 const std::shared_ptr<DataProcessor>& processor,
35 const std::shared_ptr<SignalSourceContainer>& container)
36{
37 if (m_is_processing.load(std::memory_order_acquire)) {
38 m_pending_removal.emplace_back(processor, container);
39 return;
40 }
41
42 remove_processor_direct(processor, container);
43}
44
46 const std::shared_ptr<DataProcessor>& processor,
47 const std::shared_ptr<SignalSourceContainer>& container)
48{
49 auto it = m_container_processors.find(container);
50 if (it == m_container_processors.end()) {
51 return;
52 }
53
54 auto& processors = it->second;
55 auto proc_it = std::ranges::find(processors, processor);
56
57 if (proc_it != processors.end()) {
58 processor->on_detach(container);
59 processors.erase(proc_it);
60 m_processor_tags.erase(processor);
61
62 if (processors.empty()) {
63 m_container_processors.erase(it);
64 }
65 }
66}
67
69{
70 for (auto& [processor, container] : m_pending_removal) {
71 remove_processor_direct(processor, container);
72 }
73 m_pending_removal.clear();
74}
75
76void DataProcessingChain::process(const std::shared_ptr<SignalSourceContainer>& container)
77{
78 bool expected = false;
79 if (!m_is_processing.compare_exchange_strong(expected, true,
80 std::memory_order_acquire, std::memory_order_relaxed)) {
81 return;
82 }
83
84 auto it = m_container_processors.find(container);
85 if (it != m_container_processors.end()) {
86 for (auto& processor : it->second) {
87 processor->process(container);
88 }
89 }
90
91 m_is_processing.store(false, std::memory_order_release);
93}
94
96 const std::shared_ptr<SignalSourceContainer>& container,
97 const std::function<bool(const std::shared_ptr<DataProcessor>&)>& filter)
98{
99 bool expected = false;
100 if (!m_is_processing.compare_exchange_strong(expected, true,
101 std::memory_order_acquire, std::memory_order_relaxed)) {
102 return;
103 }
104
105 auto it = m_container_processors.find(container);
106 if (it != m_container_processors.end()) {
107 for (auto& processor : it->second) {
108 if (filter(processor)) {
109 processor->process(container);
110 }
111 }
112 }
113
114 m_is_processing.store(false, std::memory_order_release);
116}
117
119 const std::shared_ptr<SignalSourceContainer>& container,
120 const std::string& tag)
121{
122 bool expected = false;
123 if (!m_is_processing.compare_exchange_strong(expected, true,
124 std::memory_order_acquire, std::memory_order_relaxed)) {
125 return;
126 }
127
128 auto it = m_container_processors.find(container);
129 if (it != m_container_processors.end()) {
130 for (auto& processor : it->second) {
131 auto tag_it = m_processor_tags.find(processor);
132 if (tag_it != m_processor_tags.end() && tag_it->second == tag) {
133 processor->process(container);
134 }
135 }
136 }
137
138 m_is_processing.store(false, std::memory_order_release);
140}
141
142} // namespace MayaFlux::Kakshya
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.