forked from dotnet/reactive
/
ToObservable.cs
78 lines (65 loc) · 2.47 KB
/
ToObservable.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
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
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the Apache 2.0 License.
// See the LICENSE file in the project root for more information.
using System.Collections.Generic;
namespace System.Linq
{
public static partial class AsyncEnumerable
{
public static IObservable<TSource> ToObservable<TSource>(this IAsyncEnumerable<TSource> source)
{
if (source == null)
throw Error.ArgumentNull(nameof(source));
return new ToObservableObservable<TSource>(source);
}
private sealed class ToObservableObservable<T> : IObservable<T>
{
private readonly IAsyncEnumerable<T> _source;
public ToObservableObservable(IAsyncEnumerable<T> source)
{
_source = source;
}
public IDisposable Subscribe(IObserver<T> observer)
{
var ctd = new CancellationTokenDisposable();
async void Core()
{
await using (var e = _source.GetAsyncEnumerator(ctd.Token))
{
do
{
bool hasNext;
var value = default(T)!;
try
{
hasNext = await e.MoveNextAsync().ConfigureAwait(false);
if (hasNext)
{
value = e.Current;
}
}
catch (Exception ex)
{
if (!ctd.Token.IsCancellationRequested)
{
observer.OnError(ex);
}
return;
}
if (!hasNext)
{
observer.OnCompleted();
return;
}
observer.OnNext(value);
}
while (!ctd.Token.IsCancellationRequested);
}
}
// Fire and forget
Core();
return ctd;
}
}
}
}