MessageBroker.cs 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209
  1. using System;
  2. using System.Collections.Generic;
  3. using UniRx.InternalUtil;
  4. namespace UniRx
  5. {
  6. public interface IMessagePublisher
  7. {
  8. /// <summary>
  9. /// Send Message to all receiver.
  10. /// </summary>
  11. void Publish<T>(T message);
  12. }
  13. public interface IMessageReceiver
  14. {
  15. /// <summary>
  16. /// Subscribe typed message.
  17. /// </summary>
  18. IObservable<T> Receive<T>();
  19. }
  20. public interface IMessageBroker : IMessagePublisher, IMessageReceiver
  21. {
  22. }
  23. public interface IAsyncMessagePublisher
  24. {
  25. /// <summary>
  26. /// Send Message to all receiver and await complete.
  27. /// </summary>
  28. IObservable<Unit> PublishAsync<T>(T message);
  29. }
  30. public interface IAsyncMessageReceiver
  31. {
  32. /// <summary>
  33. /// Subscribe typed message.
  34. /// </summary>
  35. IDisposable Subscribe<T>(Func<T, IObservable<Unit>> asyncMessageReceiver);
  36. }
  37. public interface IAsyncMessageBroker : IAsyncMessagePublisher, IAsyncMessageReceiver
  38. {
  39. }
  40. /// <summary>
  41. /// In-Memory PubSub filtered by Type.
  42. /// </summary>
  43. public class MessageBroker : IMessageBroker, IDisposable
  44. {
  45. /// <summary>
  46. /// MessageBroker in Global scope.
  47. /// </summary>
  48. public static readonly IMessageBroker Default = new MessageBroker();
  49. bool isDisposed = false;
  50. readonly Dictionary<Type, object> notifiers = new Dictionary<Type, object>();
  51. public void Publish<T>(T message)
  52. {
  53. object notifier;
  54. lock (notifiers)
  55. {
  56. if (isDisposed) return;
  57. if (!notifiers.TryGetValue(typeof(T), out notifier))
  58. {
  59. return;
  60. }
  61. }
  62. ((ISubject<T>)notifier).OnNext(message);
  63. }
  64. public IObservable<T> Receive<T>()
  65. {
  66. object notifier;
  67. lock (notifiers)
  68. {
  69. if (isDisposed) throw new ObjectDisposedException("MessageBroker");
  70. if (!notifiers.TryGetValue(typeof(T), out notifier))
  71. {
  72. ISubject<T> n = new Subject<T>().Synchronize();
  73. notifier = n;
  74. notifiers.Add(typeof(T), notifier);
  75. }
  76. }
  77. return ((IObservable<T>)notifier).AsObservable();
  78. }
  79. public void Dispose()
  80. {
  81. lock (notifiers)
  82. {
  83. if (!isDisposed)
  84. {
  85. isDisposed = true;
  86. notifiers.Clear();
  87. }
  88. }
  89. }
  90. }
  91. /// <summary>
  92. /// In-Memory PubSub filtered by Type.
  93. /// </summary>
  94. public class AsyncMessageBroker : IAsyncMessageBroker, IDisposable
  95. {
  96. /// <summary>
  97. /// AsyncMessageBroker in Global scope.
  98. /// </summary>
  99. public static readonly IAsyncMessageBroker Default = new AsyncMessageBroker();
  100. bool isDisposed = false;
  101. readonly Dictionary<Type, object> notifiers = new Dictionary<Type, object>();
  102. public IObservable<Unit> PublishAsync<T>(T message)
  103. {
  104. UniRx.InternalUtil.ImmutableList<Func<T, IObservable<Unit>>> notifier;
  105. lock (notifiers)
  106. {
  107. if (isDisposed) throw new ObjectDisposedException("AsyncMessageBroker");
  108. object _notifier;
  109. if (notifiers.TryGetValue(typeof(T), out _notifier))
  110. {
  111. notifier = (UniRx.InternalUtil.ImmutableList<Func<T, IObservable<Unit>>>)_notifier;
  112. }
  113. else
  114. {
  115. return Observable.ReturnUnit();
  116. }
  117. }
  118. var data = notifier.Data;
  119. var awaiter = new IObservable<Unit>[data.Length];
  120. for (int i = 0; i < data.Length; i++)
  121. {
  122. awaiter[i] = data[i].Invoke(message);
  123. }
  124. return Observable.WhenAll(awaiter);
  125. }
  126. public IDisposable Subscribe<T>(Func<T, IObservable<Unit>> asyncMessageReceiver)
  127. {
  128. lock (notifiers)
  129. {
  130. if (isDisposed) throw new ObjectDisposedException("AsyncMessageBroker");
  131. object _notifier;
  132. if (!notifiers.TryGetValue(typeof(T), out _notifier))
  133. {
  134. var notifier = UniRx.InternalUtil.ImmutableList<Func<T, IObservable<Unit>>>.Empty;
  135. notifier = notifier.Add(asyncMessageReceiver);
  136. notifiers.Add(typeof(T), notifier);
  137. }
  138. else
  139. {
  140. var notifier = (ImmutableList<Func<T, IObservable<Unit>>>)_notifier;
  141. notifier = notifier.Add(asyncMessageReceiver);
  142. notifiers[typeof(T)] = notifier;
  143. }
  144. }
  145. return new Subscription<T>(this, asyncMessageReceiver);
  146. }
  147. public void Dispose()
  148. {
  149. lock (notifiers)
  150. {
  151. if (!isDisposed)
  152. {
  153. isDisposed = true;
  154. notifiers.Clear();
  155. }
  156. }
  157. }
  158. class Subscription<T> : IDisposable
  159. {
  160. readonly AsyncMessageBroker parent;
  161. readonly Func<T, IObservable<Unit>> asyncMessageReceiver;
  162. public Subscription(AsyncMessageBroker parent, Func<T, IObservable<Unit>> asyncMessageReceiver)
  163. {
  164. this.parent = parent;
  165. this.asyncMessageReceiver = asyncMessageReceiver;
  166. }
  167. public void Dispose()
  168. {
  169. lock (parent.notifiers)
  170. {
  171. object _notifier;
  172. if (parent.notifiers.TryGetValue(typeof(T), out _notifier))
  173. {
  174. var notifier = (ImmutableList<Func<T, IObservable<Unit>>>)_notifier;
  175. notifier = notifier.Remove(asyncMessageReceiver);
  176. parent.notifiers[typeof(T)] = notifier;
  177. }
  178. }
  179. }
  180. }
  181. }
  182. }