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(); } } } } }