Files
Continentis/Assets/Plugins/UniRx/Scripts/Subjects/Subject.cs

148 lines
4.2 KiB
C#
Raw Normal View History

2025-10-03 00:02:43 -04:00
using System;
using UniRx.InternalUtil;
namespace UniRx
{
public sealed class Subject<T> : ISubject<T>, IDisposable, IOptimizedObservable<T>
{
2026-03-20 11:56:50 -04:00
private bool isDisposed;
2025-10-03 00:02:43 -04:00
2026-03-20 11:56:50 -04:00
private bool isStopped;
private Exception lastError;
private readonly object observerLock = new();
private IObserver<T> outObserver = EmptyObserver<T>.Instance;
2025-10-03 00:02:43 -04:00
2026-03-20 11:56:50 -04:00
public bool HasObservers => !(outObserver is EmptyObserver<T>) && !isStopped && !isDisposed;
public void Dispose()
2025-10-03 00:02:43 -04:00
{
2026-03-20 11:56:50 -04:00
lock (observerLock)
2025-10-03 00:02:43 -04:00
{
2026-03-20 11:56:50 -04:00
isDisposed = true;
outObserver = DisposedObserver<T>.Instance;
2025-10-03 00:02:43 -04:00
}
}
2026-03-20 11:56:50 -04:00
public bool IsRequiredSubscribeOnCurrentThread()
{
return false;
}
2025-10-03 00:02:43 -04:00
public void OnCompleted()
{
IObserver<T> old;
lock (observerLock)
{
ThrowIfDisposed();
if (isStopped) return;
old = outObserver;
outObserver = EmptyObserver<T>.Instance;
isStopped = true;
}
old.OnCompleted();
}
public void OnError(Exception error)
{
if (error == null) throw new ArgumentNullException("error");
IObserver<T> old;
lock (observerLock)
{
ThrowIfDisposed();
if (isStopped) return;
old = outObserver;
outObserver = EmptyObserver<T>.Instance;
isStopped = true;
lastError = error;
}
old.OnError(error);
}
public void OnNext(T value)
{
outObserver.OnNext(value);
}
public IDisposable Subscribe(IObserver<T> observer)
{
if (observer == null) throw new ArgumentNullException("observer");
var ex = default(Exception);
lock (observerLock)
{
ThrowIfDisposed();
if (!isStopped)
{
var listObserver = outObserver as ListObserver<T>;
if (listObserver != null)
{
outObserver = listObserver.Add(observer);
}
else
{
var current = outObserver;
if (current is EmptyObserver<T>)
outObserver = observer;
else
2026-03-20 11:56:50 -04:00
outObserver =
new ListObserver<T>(new ImmutableList<IObserver<T>>(new[] { current, observer }));
2025-10-03 00:02:43 -04:00
}
return new Subscription(this, observer);
}
ex = lastError;
}
if (ex != null)
observer.OnError(ex);
else
observer.OnCompleted();
return Disposable.Empty;
}
2026-03-20 11:56:50 -04:00
private void ThrowIfDisposed()
2025-10-03 00:02:43 -04:00
{
if (isDisposed) throw new ObjectDisposedException("");
}
2026-03-20 11:56:50 -04:00
private class Subscription : IDisposable
2025-10-03 00:02:43 -04:00
{
2026-03-20 11:56:50 -04:00
private readonly object gate = new();
private Subject<T> parent;
private IObserver<T> unsubscribeTarget;
2025-10-03 00:02:43 -04:00
public Subscription(Subject<T> parent, IObserver<T> unsubscribeTarget)
{
this.parent = parent;
this.unsubscribeTarget = unsubscribeTarget;
}
public void Dispose()
{
lock (gate)
{
if (parent != null)
lock (parent.observerLock)
{
var listObserver = parent.outObserver as ListObserver<T>;
if (listObserver != null)
parent.outObserver = listObserver.Remove(unsubscribeTarget);
else
parent.outObserver = EmptyObserver<T>.Instance;
unsubscribeTarget = null;
parent = null;
}
}
}
}
}
}