using System;
using System.Collections.Generic;
using EcsRx.MicroRx.Subjects;
/*
* This code was taken from UniRx project by neuecc
* https://github.com/neuecc/UniRx
*/
namespace EcsRx.MicroRx.Events
{
public class MessageBroker : IMessageBroker, IDisposable
{
///
/// MessageBroker in Global scope.
///
public static readonly IMessageBroker Default = new MessageBroker();
bool isDisposed;
readonly Dictionary notifiers = new Dictionary();
public void Publish(T message)
{
object notifier;
lock (notifiers)
{
if (isDisposed) return;
if (!notifiers.TryGetValue(typeof(T), out notifier))
{
return;
}
}
((ISubject)notifier).OnNext(message);
}
public IObservable Receive()
{
object notifier;
lock (notifiers)
{
if (isDisposed) throw new ObjectDisposedException("MessageBroker");
if (!notifiers.TryGetValue(typeof(T), out notifier))
{
var n = new Subject();
notifier = n;
notifiers.Add(typeof(T), notifier);
}
}
return ((IObservable)notifier);
}
public void Dispose()
{
lock (notifiers)
{
if (!isDisposed)
{
isDisposed = true;
notifiers.Clear();
}
}
}
}
}