59 x.
type = Type::String;
80 int64_t toInt()
const;
81 double toFloat()
const;
83 std::string toString()
const;
86 bool equals(
const Value& o)
const;
103 other.dispose_ =
nullptr;
104 other.disposed_ =
true;
108 if (
this != &other) {
111 dispose_ = std::move(other.dispose_);
112 disposed_ = other.disposed_;
113 other.dispose_ =
nullptr;
114 other.disposed_ =
true;
123 if (dispose_ && !disposed_) {
125 auto fn = std::move(dispose_);
135 std::function<void()> dispose_;
136 bool disposed_ =
false;
148 using NextFn = std::function<void(
const T&)>;
149 using ErrorFn = std::function<void(
const std::string&)>;
169 void error(
const std::string& e)
const {
170 if (stopped_)
return;
176 if (stopped_)
return;
182 mutable bool stopped_ =
false;
202 obs.
onNext = std::move(next);
211 obs.
onNext = std::move(next);
221 template <
typename R>
244 : fn_(
std::move(
fn)) {}
259 in.
onNext = [out, pred](
const T&
v) {
262 in.
onError = [out](
const std::string& e) { out.
error(e); };
264 return self->subscribe(std::move(in));
276 in.
onError = [out](
const std::string& e) { out.
error(e); };
278 return self->subscribe(std::move(in));
287 auto remaining = std::make_shared<int>(
n);
289 in.
onNext = [out, remaining](
const T&
v) {
290 if (*remaining <= 0)
return;
293 if (*remaining == 0) {
298 in.
onError = [out](
const std::string& e) { out.
error(e); };
300 return self->subscribe(std::move(in));
309 auto remaining = std::make_shared<int>(
n);
311 in.
onNext = [out, remaining](
const T&
v) {
312 if (*remaining > 0) {
318 in.
onError = [out](
const std::string& e) { out.
error(e); };
320 return self->subscribe(std::move(in));
329 auto done = std::make_shared<bool>(
false);
331 in.
onNext = [out, done](
const T&
v) {
337 in.
onError = [out](
const std::string& e) { out.
error(e); };
339 return self->subscribe(std::move(in));
350 in.
onError = [out](
const std::string& e) { out.
error(e); };
352 auto src = std::make_shared<Subscription>(self->subscribe(std::move(in)));
355 stopper.
onNext = [out, src](
const T&)
mutable {
363 auto stop = std::make_shared<Subscription>(other->
subscribe(std::move(stopper)));
378 auto last = std::make_shared<T>();
379 auto has = std::make_shared<bool>(
false);
381 in.
onNext = [out, last, has](
const T&
v) {
382 if (*has && *last ==
v)
return;
387 in.
onError = [out](
const std::string& e) { out.
error(e); };
389 return self->subscribe(std::move(in));
407 auto slot = std::make_shared<Slot>();
408 slot->observer = std::move(obs);
411 std::lock_guard<std::mutex> lock(mu_);
419 slots_.push_back(slot);
422 return Subscription([
this, slot]() { removeSlot(slot); });
445 std::lock_guard<std::mutex> lock(mu_);
446 for (
const auto&
s : slots_)
447 if (!
s->removed)
return true;
453 std::lock_guard<std::mutex> lock(mu_);
455 for (
const auto&
s : slots_)
456 if (!
s->removed)
n++;
467 std::vector<std::shared_ptr<Slot>> copy;
469 std::lock_guard<std::mutex> lock(mu_);
470 for (
auto&
s : slots_)
473 for (
auto&
s : copy) {
474 Observer<T> snapshot =
s->observer;
479 void removeSlot(
const std::shared_ptr<Slot>& slot) {
480 std::lock_guard<std::mutex> lock(mu_);
481 slot->removed =
true;
485 std::lock_guard<std::mutex> lock(mu_);
490 mutable std::mutex mu_;
491 std::vector<std::shared_ptr<Slot>> slots_;
492 bool stopped_ =
false;
513 std::lock_guard<std::mutex> lock(mu_);
519 if (replay) obs.
next(latest);
520 return base_.subscribe(std::move(obs));
526 std::lock_guard<std::mutex> lock(mu_);
534 std::lock_guard<std::mutex> lock(mu_);
535 latest_ = std::move(
v);
538 base_.onNext(emitted);
546 std::lock_guard<std::mutex> lock(mu_);
555 std::lock_guard<std::mutex> lock(mu_);
564 mutable std::mutex mu_;
566 bool completed_ =
false;
584 std::vector<T> replay;
587 std::lock_guard<std::mutex> lock(mu_);
590 for (
const auto&
v : replay)
592 return base_.subscribe(std::move(obs));
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());
605 base_.onNext(emitted);
619 mutable std::mutex mu_;
620 std::vector<T> buffer_;
637 T
get()
const {
return subject_.getValue(); }
639 void set(T
v) { subject_.setValue(std::move(
v)); }
#define EVENGINE_API_PLATFORM
Thread-affine callback registry with mutation-safe dispatch.
Internal: observable built directly from a subscribe function.
Subscription subscribe(Observer< T > obs) override
Subscribe.
AnonymousObservable(std::function< Subscription(Observer< T >)> fn)
Constructs a AnonymousObservable.
Subject that replays the latest value to every new subscriber. setValue()/onNext() update the stored ...
void onCompleted()
Stops replay and delivers a terminal completion.
bool hasObservers() const
True while at least one live observer is registered.
T getValue() const
Current stored value.
void onNext(T v)
Alias of setValue().
void onError(const std::string &e)
Stops replay and delivers a terminal error.
void setValue(T v)
Stores a new value and pushes it to observers.
BehaviorSubject(T initial)
Creates a subject with an initial (replayed) value.
Subscription subscribe(Observer< T > obs) override
Replays the latest value, then subscribes the observer.
Push-based stream source. Operators return a new (caller-owned) Observable; subscribe() returns a Sub...
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.
Observable< T > * distinctUntilChanged()
Suppresses consecutive duplicate values (uses operator==).
Observable< T > * takeUntil(Observable< T > *other)
Stops the stream when other emits or completes.
Observable< R > * map(std::function< R(const T &)> fn)
Transforms each value with fn.
virtual ~Observable()=default
Releases Observable resources.
Subscription subscribe(typename Observer< T >::NextFn next)
Subscribes with a value callback only.
Observable< T > * skip(int n)
Drops the first n values.
Observable< T > * take(int n)
Emits at most the first n values, then completes.
Observable< T > * first()
Emits only the first value, then completes.
Subscription subscribe(typename Observer< T >::NextFn next, typename Observer< T >::ErrorFn error, typename Observer< T >::CompletedFn completed)
Subscribes with value/error/completed callbacks.
Callback bundle pushed to a subscriber. Once error() or completed() fires the observer is stopped and...
CompletedFn onCompleted
Completion callback (terminal).
void error(const std::string &e) const
Delivers a terminal error and stops.
std::function< void(const T &)> NextFn
void completed() const
Delivers a terminal completion and stops.
NextFn onNext
Value callback.
std::function< void(const std::string &)> ErrorFn
bool isStopped() const
True after error()/completed() (or setStopped()).
ErrorFn onError
Error callback (terminal).
std::function< void()> CompletedFn
void next(const T &v) const
Delivers a value unless stopped.
void setStopped() const
Manually marks the observer stopped (used by operators).
Observable value backed by a BehaviorSubject. get() returns the current value; set() stores and pushe...
ReactiveProperty(T initial=T())
Creates a property with an initial value.
BehaviorSubject< T > * asSubject()
Underlying behavior subject / observable view.
T get() const
Current value.
void set(T v)
Stores a new value and notifies subscribers.
Subscription subscribe(typename Observer< T >::NextFn next)
Subscribe.
Subscription subscribe(Observer< T > obs)
Subscribes with a full observer or a value callback.
Observable< T > * asObservable()
As observable.
Subject that buffers up to capacity values (0 = unlimited) and replays the buffer to every new subscr...
ReplaySubject(int capacity=0)
Creates a replaying subject with the given buffer capacity.
void onNext(T v)
Buffers and pushes a new value to observers.
Subscription subscribe(Observer< T > obs) override
Replays buffered values, then subscribes the observer.
void onCompleted()
Delivers a terminal completion to observers.
void onError(const std::string &e)
Delivers a terminal error to observers.
bool hasObservers() const
True while at least one live observer is registered.
Multicast push-based stream: both an Observable and a push source. Thread-safe: onNext/onError/onComp...
int observerCount() const
Number of live (non-disposed) observers.
~Subject() override=default
Releases Subject resources.
Subscription subscribe(Observer< T > obs) override
Registers an observer; returns a Subscription that unregisters it.
bool hasObservers() const
True while at least one live observer is registered.
void onCompleted()
Pushes a terminal completion and stops the subject.
void onNext(const T &v)
Pushes a value to every registered observer.
void onError(const std::string &e)
Pushes a terminal error and stops the subject.
RAII-style dispose handle for an active stream subscription. Disposing unsubscribes from the source; ...
void dispose()
Runs the dispose callback once; idempotent.
Subscription(std::function< void()> dispose)
Wraps a dispose callback (usually unsubscribing from a Subject).
bool isDisposed() const
True once dispose() has run (or the handle was moved from).
Subscription & operator=(Subscription &&other) noexcept
Operator =.
~Subscription()
Releases Subscription resources.
Subscription(const Subscription &)=delete
Subscription & operator=(const Subscription &)=delete
Subscription()=default
Constructs a Subscription.
Subscription(Subscription &&other) noexcept
Constructs a Subscription.
Runtime variant used by the script-facing (Squirrel) streams. The C++ core is templated (Observer<T>/...
bool isNil() const
Type predicates.
static Value makeFloat(double v)
Constructs a floating-point value.
static Value makeNil()
Constructs a nil value.
static Value makePtr(void *v)
Constructs a pointer value (borrowed, not owned).
static Value makeInt(int64_t v)
Constructs an integer value.
static Value makeString(std::string v)
Constructs a string value (takes ownership).
static Value makeBool(bool v)
Constructs a boolean value.
enum EVENGINE_API_FOUNDATION Float
OT_FLOAT.
enum EVENGINE_API_FOUNDATION Bool
OT_BOOL.
enum EVENGINE_API_FOUNDATION String
OT_STRING.
SettlementPipeline::Stage fn