forked from ravendb/ravendb
/
Observing.cs
50 lines (44 loc) · 1.03 KB
/
Observing.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Reactive.Disposables;
using Raven.Database.Util;
namespace Raven.SimulatedWorkLoad
{
public class Observing<T> : IObservable<T>
{
private readonly IEnumerator<T> enumerator;
private readonly ConcurrentSet<IObserver<T>> observers = new ConcurrentSet<IObserver<T>>();
public bool Completed { get; private set; }
public Observing(IEnumerable<T> src)
{
enumerator = src.GetEnumerator();
}
public IDisposable Subscribe(IObserver<T> observer)
{
observers.Add(observer);
return Disposable.Create(() => observers.TryRemove(observer));
}
public void Release(int count)
{
if (Completed)
return;
for (int i = 0; i < count; i++)
{
if (enumerator.MoveNext() == false)
{
foreach (var observer in observers)
{
observer.OnCompleted();
}
Completed = true;
break;
}
foreach (var observer in observers)
{
observer.OnNext(enumerator.Current);
}
}
}
}
}