The Pedigree Project 0.1
FileEvent.cc
1/*
2 * Copyright (c) 2026, Pedigree Developers
3 *
4 * Permission to use, copy, modify, and distribute this software for any
5 * purpose with or without fee is hereby granted, provided that the above
6 * copyright notice and this permission notice appear in all copies.
7 */
8
9#include "FileEvent.h"
10#include "pedigree/kernel/LockGuard.h"
11#include "pedigree/kernel/Log.h"
12#include "pedigree/kernel/process/Mutex.h"
13#include "pedigree/kernel/process/OperationBarrier.h"
14#include "pedigree/kernel/utilities/List.h"
15#include "pedigree/kernel/utilities/assert.h"
16#include "pedigree/kernel/utilities/utility.h"
17
18namespace {
19// Closed subscriptions stay counted until reset, so this is conservative.
20size_t fileEventSubscriptions = 0;
21} // namespace
22
24 public:
25 FileEventTarget(FileEventMask interest, FileEventObserver* observer)
26 : m_Interest(interest), m_pObserver(observer), m_Notifications() {}
27
29 if (!m_Notifications.isClosedAndDrained()) {
30 m_Notifications.closeAndWait();
31 }
32 }
33
34 bool interestedIn(FileEventMask mask) const {
35 return (mask & m_Interest) != 0;
36 }
37
38 FileEventMask interest() const {
39 return m_Interest;
40 }
41
42 bool admit() {
43 return m_Notifications.tryEnter();
44 }
45
46 void notifyAdmitted(const FileEvent& event) {
47 m_pObserver->fileEvent(event);
48 m_Notifications.leave();
49 }
50
51 void closeAdmission() {
52 m_Notifications.close();
53 }
54
55 void retire() {
56 m_Notifications.closeAndWait();
57 }
58
59 size_t sequence = 0;
60
61 private:
62 FileEventMask m_Interest;
63 FileEventObserver* m_pObserver;
64 OperationBarrier m_Notifications;
65};
66
68 public:
69 FileEventState() : m_Lock(), m_Targets(), m_Open(true) {}
70
72 LockGuard<Mutex> guard(m_Lock);
73 if (m_Targets.count()) {
74 FATAL("FileEventState destroyed with live subscriptions.");
75 }
76 }
77
78 bool add(const SharedPointer<FileEventTarget>& target) {
79 LockGuard<Mutex> guard(m_Lock);
80 if (!m_Open) {
81 return false;
82 }
83 if (m_NextSequence == ~size_t(0) || !m_Targets.tryPushBack(target))
84 return false;
85 target->sequence = m_NextSequence++;
86 __atomic_store_n(&m_Interest, m_Interest | target->interest(), __ATOMIC_RELEASE);
87 __atomic_add_fetch(&fileEventSubscriptions, size_t(1), __ATOMIC_RELEASE);
88 return true;
89 }
90
91 void remove(const SharedPointer<FileEventTarget>& target) {
92 {
93 LockGuard<Mutex> guard(m_Lock);
94 for (auto it = m_Targets.begin(); it != m_Targets.end(); ++it) {
95 if (*it == target) {
96 m_Targets.erase(it);
97 FileEventMask interest = 0;
98 if (m_Open) {
99 for (const auto& remaining : m_Targets) {
100 interest |= remaining->interest();
101 }
102 }
103 __atomic_store_n(&m_Interest, interest, __ATOMIC_RELEASE);
104 __atomic_sub_fetch(&fileEventSubscriptions, size_t(1), __ATOMIC_RELEASE);
105 break;
106 }
107 }
108 }
109 target->retire();
110 }
111
112 bool hasTargets(FileEventMask mask) const {
113 return (__atomic_load_n(&m_Interest, __ATOMIC_ACQUIRE) & mask) != 0;
114 }
115
116 void notify(const FileEvent& event) {
117 OperationBarrier::Lease publication;
118 size_t boundary;
119 {
120 LockGuard<Mutex> guard(m_Lock);
121 if (!m_Open || !m_Targets.count())
122 return;
123 if (!m_Publications.tryAcquire(publication))
124 return;
125 boundary = m_NextSequence - 1;
126 }
127 size_t after = 0;
128 while (true) {
130 {
131 LockGuard<Mutex> guard(m_Lock);
132 if (!m_Open)
133 return;
134 for (const auto& target : m_Targets) {
135 if (target->sequence > after && target->sequence <= boundary &&
136 target->interestedIn(event.mask) && target->admit()) {
137 selected = target;
138 after = target->sequence;
139 break;
140 }
141 }
142 }
143 if (!selected)
144 return;
145 selected->notifyAdmitted(event);
146 }
147 }
148
149 void beginClose(const FileEvent* finalEvent = nullptr) {
150 size_t boundary;
151 {
152 LockGuard<Mutex> guard(m_Lock);
153 if (!m_Open)
154 return;
155 m_Open = false;
156 __atomic_store_n(&m_Interest, FileEventMask(0), __ATOMIC_RELEASE);
157 // The closing publisher is separately counted so a concurrent drain
158 // cannot return before the final callbacks have been admitted.
159 const bool admitted = m_ClosingPublication.tryEnter();
160 assert(admitted);
161 m_ClosingPublication.close();
162 m_Publications.close();
163 boundary = m_NextSequence - 1;
164 }
165 size_t after = 0;
166 while (true) {
168 bool deliver = false;
169 {
170 LockGuard<Mutex> guard(m_Lock);
171 for (const auto& target : m_Targets) {
172 if (target->sequence > after && target->sequence <= boundary) {
173 selected = target;
174 after = target->sequence;
175 deliver = finalEvent && target->interestedIn(finalEvent->mask) && target->admit();
176 target->closeAdmission();
177 break;
178 }
179 }
180 }
181 if (!selected)
182 break;
183 if (deliver)
184 selected->notifyAdmitted(*finalEvent);
185 }
186 m_ClosingPublication.leave();
187 }
188
189 void drain() {
190 {
191 LockGuard<Mutex> guard(m_Lock);
192 if (m_Open)
193 return;
194 }
195 m_ClosingPublication.wait();
196 m_Publications.wait();
197 size_t after = 0;
198 while (true) {
200 {
201 LockGuard<Mutex> guard(m_Lock);
202 for (const auto& target : m_Targets) {
203 if (target->sequence > after) {
204 selected = target;
205 after = target->sequence;
206 break;
207 }
208 }
209 }
210 if (!selected)
211 return;
212 selected->retire();
213 }
214 }
215
216 private:
217 Mutex m_Lock;
219 bool m_Open;
220 FileEventMask m_Interest = 0;
221 size_t m_NextSequence = 1;
222 OperationBarrier m_Publications;
223 OperationBarrier m_ClosingPublication;
224};
225
226FileEventSubscription::FileEventSubscription() : m_State(), m_Target(), m_Observer() {}
227
228FileEventSubscription::FileEventSubscription(FileEventSubscription&& other) noexcept
229 : m_State(pedigree_std::move(other.m_State)),
230 m_Target(pedigree_std::move(other.m_Target)),
231 m_Observer(pedigree_std::move(other.m_Observer)) {}
232
233FileEventSubscription::~FileEventSubscription() {
234 reset();
235}
236
237FileEventSubscription& FileEventSubscription::operator=(FileEventSubscription&& other) noexcept {
238 if (this != &other) {
239 reset();
240 m_State = pedigree_std::move(other.m_State);
241 m_Target = pedigree_std::move(other.m_Target);
242 m_Observer = pedigree_std::move(other.m_Observer);
243 }
244 return *this;
245}
246
247FileEventSubscription::operator bool() const {
248 return static_cast<bool>(m_State) && static_cast<bool>(m_Target);
249}
250
251void FileEventSubscription::reset() {
252 SharedPointer<FileEventState> state = m_State;
253 SharedPointer<FileEventTarget> target = m_Target;
254 if (state && target) {
255 state->remove(target);
256 }
257 m_Target.reset();
258 m_Observer.reset();
259 m_State.reset();
260}
261
262FileEventSource::FileEventSource() : m_FileEventState(new FileEventState) {}
263
264FileEventSource::~FileEventSource() {
265 closeFileEvents();
266}
267
268bool FileEventSource::subscribeFileEvents(FileEventMask interest,
269 const SharedPointer<FileEventObserver>& observer,
270 FileEventSubscription& subscription) {
271 subscription.reset();
272 if (!interest || !observer) {
273 return false;
274 }
275
276 SharedPointer<FileEventTarget> target(new FileEventTarget(interest, observer.get()));
277 if (!target || !m_FileEventState)
278 return false;
279 subscription.m_State = m_FileEventState;
280 subscription.m_Target = target;
281 subscription.m_Observer = observer;
282 if (!m_FileEventState->add(target)) {
283 subscription.reset();
284 return false;
285 }
286 return true;
287}
288
290 if (hasFileEventObservers(event.mask)) {
291 SharedPointer<FileEventState> state = m_FileEventState;
292 state->notify(event);
293 }
294}
295
297 return __atomic_load_n(&fileEventSubscriptions, __ATOMIC_ACQUIRE) != 0;
298}
299
300bool FileEventSource::hasFileEventObservers(FileEventMask mask) const {
301 return m_FileEventState && m_FileEventState->hasTargets(mask);
302}
303
305 if (event.mask && m_FileEventState) {
306 SharedPointer<FileEventState> state = m_FileEventState;
307 state->beginClose(&event);
308 state->drain();
309 }
310}
311
312void FileEventSource::closeFileEvents() {
313 SharedPointer<FileEventState> state = m_FileEventState;
314 if (!state)
315 return;
316 state->beginClose();
317 state->drain();
318}
319
321 if (m_FileEventState)
322 m_FileEventState->beginClose(&event);
323}
324
325void FileEventSource::drainFileEvents() {
326 if (m_FileEventState)
327 m_FileEventState->drain();
328}
329
330void InodeEventSource::publish(const FileEvent& event) {
331 notifyFileEvent(event);
332}
333
334void InodeEventSource::beginRetirement() {
335 beginFinalFileEvent(FileEvent(FileEvents::SourceRetired, StringView(), false));
336}
337
339 drainFileEvents();
340}
virtual void fileEvent(const FileEvent &event)=0
void notifyFileEvent(const FileEvent &event)
Definition FileEvent.cc:289
void beginFinalFileEvent(const FileEvent &event)
Definition FileEvent.cc:320
void notifyFinalFileEvent(const FileEvent &event)
Definition FileEvent.cc:304
static bool anyFileEventObservers()
Definition FileEvent.cc:296
void finishRetirement()
Definition FileEvent.cc:338
Definition List.h:61
Iterator begin()
Definition List.h:122
Iterator end()
Definition List.h:132
Definition Mutex.h:56
MUST_USE_RESULT bool tryAcquire(Lease &lease)
MUST_USE_RESULT bool tryEnter()
T * get() const
Iterator erase(Iterator &Iter)
Definition List.h:352
bool tryPushBack(const T &value)
Definition List.h:246
size_t count() const
Definition List.h:212