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 Observαble;
#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
}
}
|