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 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159
|
// Copyright (c) Microsoft Open Technologies, Inc. All rights reserved. See License.txt in the project root for license information.
using System.Collections.Generic;
using System.Reactive.Concurrency;
using System.Reactive.Disposables;
namespace System.Reactive.Linq
{
#if !NO_PERF
using ObservableImpl;
#endif
internal partial class QueryLanguage
{
#region + Subscribe +
public virtual IDisposable Subscribe<TSource>(IEnumerable<TSource> source, IObserver<TSource> observer)
{
return Subscribe_<TSource>(source, observer, SchedulerDefaults.Iteration);
}
public virtual IDisposable Subscribe<TSource>(IEnumerable<TSource> source, IObserver<TSource> observer, IScheduler scheduler)
{
return Subscribe_<TSource>(source, observer, scheduler);
}
private static IDisposable Subscribe_<TSource>(IEnumerable<TSource> source, IObserver<TSource> observer, IScheduler scheduler)
{
#if !NO_PERF
//
// [OK] Use of unsafe Subscribe: we're calling into a known producer implementation.
//
return new ToObservable<TSource>(source, scheduler).Subscribe/*Unsafe*/(observer);
#else
var e = source.GetEnumerator();
var flag = new BooleanDisposable();
scheduler.Schedule(self =>
{
var hasNext = false;
var ex = default(Exception);
var current = default(TSource);
if (flag.IsDisposed)
{
e.Dispose();
return;
}
try
{
hasNext = e.MoveNext();
if (hasNext)
current = e.Current;
}
catch (Exception exception)
{
ex = exception;
}
if (!hasNext || ex != null)
{
e.Dispose();
}
if (ex != null)
{
observer.OnError(ex);
return;
}
if (!hasNext)
{
observer.OnCompleted();
return;
}
observer.OnNext(current);
self();
});
return flag;
#endif
}
#endregion
#region + ToEnumerable +
public virtual IEnumerable<TSource> ToEnumerable<TSource>(IObservable<TSource> source)
{
return new AnonymousEnumerable<TSource>(() => source.GetEnumerator());
}
#endregion
#region ToEvent
public virtual IEventSource<Unit> ToEvent(IObservable<Unit> source)
{
return new EventSource<Unit>(source, (h, _) => h(Unit.Default));
}
public virtual IEventSource<TSource> ToEvent<TSource>(IObservable<TSource> source)
{
return new EventSource<TSource>(source, (h, value) => h(value));
}
#endregion
#region ToEventPattern
public virtual IEventPatternSource<TEventArgs> ToEventPattern<TEventArgs>(IObservable<EventPattern<TEventArgs>> source)
#if !NO_EVENTARGS_CONSTRAINT
where TEventArgs : EventArgs
#endif
{
return new EventPatternSource<TEventArgs>(
#if !NO_VARIANCE
source,
#else
source.Select(x => (EventPattern<object, TEventArgs>)x),
#endif
(h, evt) => h(evt.Sender, evt.EventArgs)
);
}
#endregion
#region + ToObservable +
public virtual IObservable<TSource> ToObservable<TSource>(IEnumerable<TSource> source)
{
#if !NO_PERF
return new ToObservable<TSource>(source, SchedulerDefaults.Iteration);
#else
return ToObservable_(source, SchedulerDefaults.Iteration);
#endif
}
public virtual IObservable<TSource> ToObservable<TSource>(IEnumerable<TSource> source, IScheduler scheduler)
{
#if !NO_PERF
return new ToObservable<TSource>(source, scheduler);
#else
return ToObservable_(source, scheduler);
#endif
}
#if NO_PERF
private static IObservable<TSource> ToObservable_<TSource>(IEnumerable<TSource> source, IScheduler scheduler)
{
return new AnonymousObservable<TSource>(observer => source.Subscribe(observer, scheduler));
}
#endif
#endregion
}
}
|