Collect.cs 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139
  1. // Copyright (c) Microsoft Open Technologies, Inc. All rights reserved. See License.txt in the project root for license information.
  2. #if !NO_PERF
  3. using System;
  4. using System.Collections.Generic;
  5. using System.Reactive;
  6. using System.Reactive.Threading;
  7. using System.Threading;
  8. namespace System.Reactive.Linq.ObservableImpl
  9. {
  10. class Collect<TSource, TResult> : PushToPullAdapter<TSource, TResult>
  11. {
  12. private readonly Func<TResult> _getInitialCollector;
  13. private readonly Func<TResult, TSource, TResult> _merge;
  14. private readonly Func<TResult, TResult> _getNewCollector;
  15. public Collect(IObservable<TSource> source, Func<TResult> getInitialCollector, Func<TResult, TSource, TResult> merge, Func<TResult, TResult> getNewCollector)
  16. : base(source)
  17. {
  18. _getInitialCollector = getInitialCollector;
  19. _merge = merge;
  20. _getNewCollector = getNewCollector;
  21. }
  22. protected override PushToPullSink<TSource, TResult> Run(IDisposable subscription)
  23. {
  24. var sink = new _(this, subscription);
  25. sink.Initialize();
  26. return sink;
  27. }
  28. class _ : PushToPullSink<TSource, TResult>
  29. {
  30. private readonly Collect<TSource, TResult> _parent;
  31. public _(Collect<TSource, TResult> parent, IDisposable subscription)
  32. : base(subscription)
  33. {
  34. _parent = parent;
  35. }
  36. private object _gate;
  37. private TResult _collector;
  38. private bool _hasFailed;
  39. private Exception _error;
  40. private bool _hasCompleted;
  41. private bool _done;
  42. public void Initialize()
  43. {
  44. _gate = new object();
  45. _collector = _parent._getInitialCollector();
  46. }
  47. public override void OnNext(TSource value)
  48. {
  49. lock (_gate)
  50. {
  51. try
  52. {
  53. _collector = _parent._merge(_collector, value);
  54. }
  55. catch (Exception ex)
  56. {
  57. _error = ex;
  58. _hasFailed = true;
  59. base.Dispose();
  60. }
  61. }
  62. }
  63. public override void OnError(Exception error)
  64. {
  65. base.Dispose();
  66. lock (_gate)
  67. {
  68. _error = error;
  69. _hasFailed = true;
  70. }
  71. }
  72. public override void OnCompleted()
  73. {
  74. base.Dispose();
  75. lock (_gate)
  76. {
  77. _hasCompleted = true;
  78. }
  79. }
  80. public override bool TryMoveNext(out TResult current)
  81. {
  82. lock (_gate)
  83. {
  84. if (_hasFailed)
  85. {
  86. current = default(TResult);
  87. _error.Throw();
  88. }
  89. else
  90. {
  91. if (_hasCompleted)
  92. {
  93. if (_done)
  94. {
  95. current = default(TResult);
  96. return false;
  97. }
  98. current = _collector;
  99. _done = true;
  100. }
  101. else
  102. {
  103. current = _collector;
  104. try
  105. {
  106. _collector = _parent._getNewCollector(current);
  107. }
  108. catch
  109. {
  110. base.Dispose();
  111. throw;
  112. }
  113. }
  114. }
  115. return true;
  116. }
  117. }
  118. }
  119. }
  120. }
  121. #endif