载入中...
搜索中...
未找到
Rx.h
浏览该文件的文档.
1#pragma once
2#include "common/Export.h"
3
4
5#include <functional>
6#include <memory>
7#include <mutex>
8#include <string>
9#include <vector>
10
11#include "common/Exception.h"
12
13namespace eve::rx {
14
22public:
24 enum class Type { Nil, Int, Float, Bool, String, Ptr };
25
26 Type type = Type::Nil;
27 int64_t i = 0;
28 double f = 0;
29 bool b = false;
30 std::string s;
31 void* p = nullptr;
32
34 static Value makeNil() { return {}; }
36 static Value makeInt(int64_t v) {
37 Value x;
38 x.type = Type::Int;
39 x.i = v;
40 return x;
41 }
43 static Value makeFloat(double v) {
44 Value x;
45 x.type = Type::Float;
46 x.f = v;
47 return x;
48 }
50 static Value makeBool(bool v) {
51 Value x;
52 x.type = Type::Bool;
53 x.b = v;
54 return x;
55 }
57 static Value makeString(std::string v) {
58 Value x;
59 x.type = Type::String;
60 x.s = std::move(v);
61 return x;
62 }
64 static Value makePtr(void* v) {
65 Value x;
66 x.type = Type::Ptr;
67 x.p = v;
68 return x;
69 }
70
72 bool isNil() const { return type == Type::Nil; }
73 bool isInt() const { return type == Type::Int; }
74 bool isFloat() const { return type == Type::Float; }
75 bool isBool() const { return type == Type::Bool; }
76 bool isString() const { return type == Type::String; }
77 bool isPtr() const { return type == Type::Ptr; }
78
80 int64_t toInt() const;
81 double toFloat() const;
82 bool toBool() const;
83 std::string toString() const;
84
86 bool equals(const Value& o) const;
87};
88
94public:
96 Subscription() = default;
98 explicit Subscription(std::function<void()> dispose) : dispose_(std::move(dispose)) {}
99 Subscription(const Subscription&) = delete;
102 Subscription(Subscription&& other) noexcept : dispose_(std::move(other.dispose_)), disposed_(other.disposed_) {
103 other.dispose_ = nullptr;
104 other.disposed_ = true;
105 }
108 if (this != &other) {
110 dispose();
111 dispose_ = std::move(other.dispose_);
112 disposed_ = other.disposed_;
113 other.dispose_ = nullptr;
114 other.disposed_ = true;
115 }
116 return *this;
117 }
120
122 void dispose() {
123 if (dispose_ && !disposed_) {
124 disposed_ = true;
125 auto fn = std::move(dispose_);
126 dispose_ = nullptr;
128 fn();
129 }
130 }
132 bool isDisposed() const { return disposed_; }
133
134private:
135 std::function<void()> dispose_;
136 bool disposed_ = false;
137};
138
144template <typename T>
146class Observer {
147public:
148 using NextFn = std::function<void(const T&)>;
149 using ErrorFn = std::function<void(const std::string&)>;
150 using CompletedFn = std::function<void()>;
151
158
160 bool isStopped() const { return stopped_; }
162 void setStopped() const { stopped_ = true; }
163
165 void next(const T& v) const {
166 if (!stopped_ && onNext) onNext(v);
167 }
169 void error(const std::string& e) const {
170 if (stopped_) return;
171 stopped_ = true;
172 if (onError) onError(e);
173 }
175 void completed() const {
176 if (stopped_) return;
177 stopped_ = true;
179 }
180
181private:
182 mutable bool stopped_ = false;
183};
184
189template <typename T>
192public:
194 virtual ~Observable() = default;
195
198
201 Observer<T> obs;
202 obs.onNext = std::move(next);
204 return subscribe(std::move(obs));
205 }
209 typename Observer<T>::CompletedFn completed) {
210 Observer<T> obs;
211 obs.onNext = std::move(next);
212 obs.onError = std::move(error);
213 obs.onCompleted = std::move(completed);
215 return subscribe(std::move(obs));
216 }
217
219 Observable<T>* filter(std::function<bool(const T&)> pred);
221 template <typename R>
223 Observable<R>* map(std::function<R(const T&)> fn);
234};
235
237template <typename T>
240public:
244 : fn_(std::move(fn)) {}
246 Subscription subscribe(Observer<T> obs) override { return fn_(std::move(obs)); }
247
248private:
249 std::function<Subscription(Observer<T>)> fn_;
250};
251
252// ---- Operators ----
253template <typename T>
255Observable<T>* Observable<T>::filter(std::function<bool(const T&)> pred) {
256 auto* self = this;
257 return new AnonymousObservable<T>([self, pred](Observer<T> out) {
258 Observer<T> in;
259 in.onNext = [out, pred](const T& v) {
260 if (pred(v)) out.next(v);
261 };
262 in.onError = [out](const std::string& e) { out.error(e); };
263 in.onCompleted = [out]() { out.completed(); };
264 return self->subscribe(std::move(in));
265 });
266}
267
268template <typename T>
269template <typename R>
271Observable<R>* Observable<T>::map(std::function<R(const T&)> fn) {
272 auto* self = this;
273 return new AnonymousObservable<R>([self, fn](Observer<R> out) {
274 Observer<T> in;
275 in.onNext = [out, fn](const T& v) { out.next(fn(v)); };
276 in.onError = [out](const std::string& e) { out.error(e); };
277 in.onCompleted = [out]() { out.completed(); };
278 return self->subscribe(std::move(in));
279 });
280}
281
282template <typename T>
285 auto* self = this;
286 return new AnonymousObservable<T>([self, n](Observer<T> out) {
287 auto remaining = std::make_shared<int>(n);
288 Observer<T> in;
289 in.onNext = [out, remaining](const T& v) {
290 if (*remaining <= 0) return;
291 *remaining -= 1;
292 out.next(v);
293 if (*remaining == 0) {
294 out.completed();
295 // note: upstream keeps running; completed() already stopped `out`
296 }
297 };
298 in.onError = [out](const std::string& e) { out.error(e); };
299 in.onCompleted = [out]() { out.completed(); };
300 return self->subscribe(std::move(in));
301 });
302}
303
304template <typename T>
307 auto* self = this;
308 return new AnonymousObservable<T>([self, n](Observer<T> out) {
309 auto remaining = std::make_shared<int>(n);
310 Observer<T> in;
311 in.onNext = [out, remaining](const T& v) {
312 if (*remaining > 0) {
313 *remaining -= 1;
314 return;
315 }
316 out.next(v);
317 };
318 in.onError = [out](const std::string& e) { out.error(e); };
319 in.onCompleted = [out]() { out.completed(); };
320 return self->subscribe(std::move(in));
321 });
322}
323
324template <typename T>
327 auto* self = this;
328 return new AnonymousObservable<T>([self](Observer<T> out) {
329 auto done = std::make_shared<bool>(false);
330 Observer<T> in;
331 in.onNext = [out, done](const T& v) {
332 if (*done) return;
333 *done = true;
334 out.next(v);
335 out.completed();
336 };
337 in.onError = [out](const std::string& e) { out.error(e); };
338 in.onCompleted = [out]() { out.completed(); };
339 return self->subscribe(std::move(in));
340 });
341}
342
343template <typename T>
346 auto* self = this;
347 return new AnonymousObservable<T>([self, other](Observer<T> out) {
348 Observer<T> in;
349 in.onNext = [out](const T& v) { out.next(v); };
350 in.onError = [out](const std::string& e) { out.error(e); };
351 in.onCompleted = [out]() { out.completed(); };
352 auto src = std::make_shared<Subscription>(self->subscribe(std::move(in)));
353
354 Observer<T> stopper;
355 stopper.onNext = [out, src](const T&) mutable {
356 out.completed();
357 src->dispose();
358 };
359 stopper.onCompleted = [out, src]() mutable {
360 out.completed();
361 src->dispose();
362 };
363 auto stop = std::make_shared<Subscription>(other->subscribe(std::move(stopper)));
364
366 return Subscription([src, stop]() mutable {
367 src->dispose();
368 stop->dispose();
369 });
370 });
371}
372
373template <typename T>
376 auto* self = this;
377 return new AnonymousObservable<T>([self](Observer<T> out) {
378 auto last = std::make_shared<T>();
379 auto has = std::make_shared<bool>(false);
380 Observer<T> in;
381 in.onNext = [out, last, has](const T& v) {
382 if (*has && *last == v) return;
383 *has = true;
384 *last = v;
385 out.next(v);
386 };
387 in.onError = [out](const std::string& e) { out.error(e); };
388 in.onCompleted = [out]() { out.completed(); };
389 return self->subscribe(std::move(in));
390 });
391}
392
397template <typename T>
399class Subject : public Observable<T> {
400public:
401 using Observable<T>::subscribe;
403 ~Subject() override = default;
404
407 auto slot = std::make_shared<Slot>();
408 slot->observer = std::move(obs);
409 {
411 std::lock_guard<std::mutex> lock(mu_);
412 if (stopped_) {
413 // Terminal state: deliver the terminal notification immediately
414 // to a fresh copy of the observer and never register the slot.
415 Observer<T> snapshot = slot->observer;
416 snapshot.completed();
417 return Subscription{};
418 }
419 slots_.push_back(slot);
420 }
422 return Subscription([this, slot]() { removeSlot(slot); });
423 }
424
426 void onNext(const T& v) { emit([&](Observer<T>& o) { o.next(v); }); }
428 void onError(const std::string& e) {
430 emit([&](Observer<T>& o) { o.error(e); });
432 markStopped();
433 }
435 void onCompleted() {
437 emit([&](Observer<T>& o) { o.completed(); });
439 markStopped();
440 }
441
443 bool hasObservers() const {
445 std::lock_guard<std::mutex> lock(mu_);
446 for (const auto& s : slots_)
447 if (!s->removed) return true;
448 return false;
449 }
451 int observerCount() const {
453 std::lock_guard<std::mutex> lock(mu_);
454 int n = 0;
455 for (const auto& s : slots_)
456 if (!s->removed) n++;
457 return n;
458 }
459
460private:
461 struct Slot {
463 bool removed = false;
464 };
465
466 void emit(const std::function<void(Observer<T>&)>& fn) {
467 std::vector<std::shared_ptr<Slot>> copy;
468 {
469 std::lock_guard<std::mutex> lock(mu_);
470 for (auto& s : slots_)
471 if (!s->removed) copy.push_back(s);
472 }
473 for (auto& s : copy) {
474 Observer<T> snapshot = s->observer;
475 fn(snapshot);
476 }
477 }
478
479 void removeSlot(const std::shared_ptr<Slot>& slot) {
480 std::lock_guard<std::mutex> lock(mu_);
481 slot->removed = true;
482 }
483
484 void markStopped() {
485 std::lock_guard<std::mutex> lock(mu_);
486 stopped_ = true;
487 slots_.clear();
488 }
489
490 mutable std::mutex mu_;
491 std::vector<std::shared_ptr<Slot>> slots_;
492 bool stopped_ = false;
493};
494
499template <typename T>
501class BehaviorSubject : public Observable<T> {
502public:
503 using Observable<T>::subscribe;
505 explicit BehaviorSubject(T initial) : latest_(std::move(initial)) {}
506
509 T latest;
510 bool replay = false;
511 {
513 std::lock_guard<std::mutex> lock(mu_);
514 if (!completed_) {
515 latest = latest_;
516 replay = true;
517 }
518 }
519 if (replay) obs.next(latest);
520 return base_.subscribe(std::move(obs));
521 }
522
524 T getValue() const {
526 std::lock_guard<std::mutex> lock(mu_);
527 return latest_;
528 }
530 void setValue(T v) {
531 T emitted;
532 {
534 std::lock_guard<std::mutex> lock(mu_);
535 latest_ = std::move(v);
536 emitted = latest_;
537 }
538 base_.onNext(emitted);
539 }
541 void onNext(T v) { setValue(std::move(v)); }
543 void onError(const std::string& e) {
544 {
546 std::lock_guard<std::mutex> lock(mu_);
547 completed_ = true;
548 }
549 base_.onError(e);
550 }
552 void onCompleted() {
553 {
555 std::lock_guard<std::mutex> lock(mu_);
556 completed_ = true;
557 }
558 base_.onCompleted();
559 }
561 bool hasObservers() const { return base_.hasObservers(); }
562
563private:
564 mutable std::mutex mu_;
565 T latest_;
566 bool completed_ = false;
567 Subject<T> base_;
568};
569
574template <typename T>
576class ReplaySubject : public Observable<T> {
577public:
578 using Observable<T>::subscribe;
580 explicit ReplaySubject(int capacity = 0) : capacity_(capacity) {}
581
584 std::vector<T> replay;
585 {
587 std::lock_guard<std::mutex> lock(mu_);
588 replay = buffer_;
589 }
590 for (const auto& v : replay)
591 if (!obs.isStopped()) obs.next(v);
592 return base_.subscribe(std::move(obs));
593 }
594
596 void onNext(T v) {
597 T emitted = v;
598 {
600 std::lock_guard<std::mutex> lock(mu_);
601 buffer_.push_back(std::move(v));
602 if (capacity_ > 0 && static_cast<int>(buffer_.size()) > capacity_)
603 buffer_.erase(buffer_.begin());
604 }
605 base_.onNext(emitted);
606 }
608 void onError(const std::string& e) {
609 base_.onError(e);
610 }
612 void onCompleted() {
613 base_.onCompleted();
614 }
616 bool hasObservers() const { return base_.hasObservers(); }
617
618private:
619 mutable std::mutex mu_;
620 std::vector<T> buffer_;
621 int capacity_;
622 Subject<T> base_;
623};
624
629template <typename T>
632public:
634 explicit ReactiveProperty(T initial = T()) : subject_(std::move(initial)) {}
635
637 T get() const { return subject_.getValue(); }
639 void set(T v) { subject_.setValue(std::move(v)); }
640
642 Subscription subscribe(Observer<T> obs) { return subject_.subscribe(std::move(obs)); }
644 Subscription subscribe(typename Observer<T>::NextFn next) { return subject_.subscribe(std::move(next)); }
645
647 BehaviorSubject<T>* asSubject() { return &subject_; }
649 Observable<T>* asObservable() { return &subject_; }
650
651private:
652 BehaviorSubject<T> subject_;
653};
654
655} // namespace eve::rx
float x
Definition AnimClip.cpp:738
int observer
Definition AnimSmr.cpp:162
const std::string & s
glm::vec4 p[6]
#define EVENGINE_API_PLATFORM
Definition Export.h:107
std::uint32_t capacity
glm::vec3 n
Definition Grass.cpp:63
float v
MeleePoint3 b
Definition MeleeHit.cpp:41
std::string error
Definition Package.cpp:60
float f
int removed
Thread-affine callback registry with mutation-safe dispatch.
Internal: observable built directly from a subscribe function.
Definition Rx.h:239
Subscription subscribe(Observer< T > obs) override
Subscribe.
Definition Rx.h:246
AnonymousObservable(std::function< Subscription(Observer< T >)> fn)
Constructs a AnonymousObservable.
Definition Rx.h:242
Subject that replays the latest value to every new subscriber. setValue()/onNext() update the stored ...
Definition Rx.h:501
void onCompleted()
Stops replay and delivers a terminal completion.
Definition Rx.h:552
bool hasObservers() const
True while at least one live observer is registered.
Definition Rx.h:561
T getValue() const
Current stored value.
Definition Rx.h:524
void onNext(T v)
Alias of setValue().
Definition Rx.h:541
void onError(const std::string &e)
Stops replay and delivers a terminal error.
Definition Rx.h:543
void setValue(T v)
Stores a new value and pushes it to observers.
Definition Rx.h:530
BehaviorSubject(T initial)
Creates a subject with an initial (replayed) value.
Definition Rx.h:505
Subscription subscribe(Observer< T > obs) override
Replays the latest value, then subscribes the observer.
Definition Rx.h:508
Push-based stream source. Operators return a new (caller-owned) Observable; subscribe() returns a Sub...
Definition Rx.h:191
virtual Subscription subscribe(Observer< T > obs)=0
Subscribes with a full observer; returns a cancel handle.
Observable< T > * filter(std::function< bool(const T &)> pred)
Passes values through only when pred(v) is true.
Definition Rx.h:255
Observable< T > * distinctUntilChanged()
Suppresses consecutive duplicate values (uses operator==).
Definition Rx.h:375
Observable< T > * takeUntil(Observable< T > *other)
Stops the stream when other emits or completes.
Definition Rx.h:345
Observable< R > * map(std::function< R(const T &)> fn)
Transforms each value with fn.
Definition Rx.h:271
virtual ~Observable()=default
Releases Observable resources.
Subscription subscribe(typename Observer< T >::NextFn next)
Subscribes with a value callback only.
Definition Rx.h:200
Observable< T > * skip(int n)
Drops the first n values.
Definition Rx.h:306
Observable< T > * take(int n)
Emits at most the first n values, then completes.
Definition Rx.h:284
Observable< T > * first()
Emits only the first value, then completes.
Definition Rx.h:326
Subscription subscribe(typename Observer< T >::NextFn next, typename Observer< T >::ErrorFn error, typename Observer< T >::CompletedFn completed)
Subscribes with value/error/completed callbacks.
Definition Rx.h:207
Callback bundle pushed to a subscriber. Once error() or completed() fires the observer is stopped and...
Definition Rx.h:146
CompletedFn onCompleted
Completion callback (terminal).
Definition Rx.h:157
void error(const std::string &e) const
Delivers a terminal error and stops.
Definition Rx.h:169
std::function< void(const T &)> NextFn
Definition Rx.h:148
void completed() const
Delivers a terminal completion and stops.
Definition Rx.h:175
NextFn onNext
Value callback.
Definition Rx.h:153
std::function< void(const std::string &)> ErrorFn
Definition Rx.h:149
bool isStopped() const
True after error()/completed() (or setStopped()).
Definition Rx.h:160
ErrorFn onError
Error callback (terminal).
Definition Rx.h:155
std::function< void()> CompletedFn
Definition Rx.h:150
void next(const T &v) const
Delivers a value unless stopped.
Definition Rx.h:165
void setStopped() const
Manually marks the observer stopped (used by operators).
Definition Rx.h:162
Observable value backed by a BehaviorSubject. get() returns the current value; set() stores and pushe...
Definition Rx.h:631
ReactiveProperty(T initial=T())
Creates a property with an initial value.
Definition Rx.h:634
BehaviorSubject< T > * asSubject()
Underlying behavior subject / observable view.
Definition Rx.h:647
T get() const
Current value.
Definition Rx.h:637
void set(T v)
Stores a new value and notifies subscribers.
Definition Rx.h:639
Subscription subscribe(typename Observer< T >::NextFn next)
Subscribe.
Definition Rx.h:644
Subscription subscribe(Observer< T > obs)
Subscribes with a full observer or a value callback.
Definition Rx.h:642
Observable< T > * asObservable()
As observable.
Definition Rx.h:649
Subject that buffers up to capacity values (0 = unlimited) and replays the buffer to every new subscr...
Definition Rx.h:576
ReplaySubject(int capacity=0)
Creates a replaying subject with the given buffer capacity.
Definition Rx.h:580
void onNext(T v)
Buffers and pushes a new value to observers.
Definition Rx.h:596
Subscription subscribe(Observer< T > obs) override
Replays buffered values, then subscribes the observer.
Definition Rx.h:583
void onCompleted()
Delivers a terminal completion to observers.
Definition Rx.h:612
void onError(const std::string &e)
Delivers a terminal error to observers.
Definition Rx.h:608
bool hasObservers() const
True while at least one live observer is registered.
Definition Rx.h:616
Multicast push-based stream: both an Observable and a push source. Thread-safe: onNext/onError/onComp...
Definition Rx.h:399
int observerCount() const
Number of live (non-disposed) observers.
Definition Rx.h:451
~Subject() override=default
Releases Subject resources.
Subscription subscribe(Observer< T > obs) override
Registers an observer; returns a Subscription that unregisters it.
Definition Rx.h:406
bool hasObservers() const
True while at least one live observer is registered.
Definition Rx.h:443
void onCompleted()
Pushes a terminal completion and stops the subject.
Definition Rx.h:435
void onNext(const T &v)
Pushes a value to every registered observer.
Definition Rx.h:426
void onError(const std::string &e)
Pushes a terminal error and stops the subject.
Definition Rx.h:428
RAII-style dispose handle for an active stream subscription. Disposing unsubscribes from the source; ...
Definition Rx.h:93
void dispose()
Runs the dispose callback once; idempotent.
Definition Rx.h:122
Subscription(std::function< void()> dispose)
Wraps a dispose callback (usually unsubscribing from a Subject).
Definition Rx.h:98
bool isDisposed() const
True once dispose() has run (or the handle was moved from).
Definition Rx.h:132
Subscription & operator=(Subscription &&other) noexcept
Operator =.
Definition Rx.h:107
~Subscription()
Releases Subscription resources.
Definition Rx.h:119
Subscription(const Subscription &)=delete
Subscription & operator=(const Subscription &)=delete
Subscription()=default
Constructs a Subscription.
Subscription(Subscription &&other) noexcept
Constructs a Subscription.
Definition Rx.h:102
Runtime variant used by the script-facing (Squirrel) streams. The C++ core is templated (Observer<T>/...
Definition Rx.h:21
Type type
Definition Rx.h:26
bool isNil() const
Type predicates.
Definition Rx.h:72
bool isInt() const
Definition Rx.h:73
static Value makeFloat(double v)
Constructs a floating-point value.
Definition Rx.h:43
bool isFloat() const
Definition Rx.h:74
static Value makeNil()
Constructs a nil value.
Definition Rx.h:34
bool isString() const
Definition Rx.h:76
static Value makePtr(void *v)
Constructs a pointer value (borrowed, not owned).
Definition Rx.h:64
std::string s
Definition Rx.h:30
Type
Type public API.
Definition Rx.h:24
static Value makeInt(int64_t v)
Constructs an integer value.
Definition Rx.h:36
static Value makeString(std::string v)
Constructs a string value (takes ownership).
Definition Rx.h:57
bool isBool() const
Definition Rx.h:75
bool isPtr() const
Definition Rx.h:77
static Value makeBool(bool v)
Constructs a boolean value.
Definition Rx.h:50
Definition Rx.cpp:11
enum EVENGINE_API_FOUNDATION Float
OT_FLOAT.
Definition Runtime.h:118
enum EVENGINE_API_FOUNDATION Bool
OT_BOOL.
Definition Runtime.h:116
enum EVENGINE_API_FOUNDATION String
OT_STRING.
Definition Runtime.h:119
SettlementPipeline::Stage fn