This repository was archived by the owner on Mar 2, 2022. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 8
/
Copy pathPublisherAsObservable.cs
122 lines (105 loc) · 2.88 KB
/
PublisherAsObservable.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
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
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
using Reactive.Streams;
using Reactor.Core;
using System.Threading;
using Reactor.Core.flow;
using Reactor.Core.subscriber;
using Reactor.Core.subscription;
using Reactor.Core.util;
namespace Reactor.Core.publisher
{
/// <summary>
/// Wraps an IPublisher and exposes it as an IObservable.
/// </summary>
/// <typeparam name="T"></typeparam>
sealed class PublisherAsObservable<T> : IObservable<T>
{
readonly IPublisher<T> source;
internal PublisherAsObservable(IPublisher<T> source)
{
this.source = source;
}
public IDisposable Subscribe(IObserver<T> observer)
{
var s = new AsObserver(observer);
source.Subscribe(s);
return s;
}
sealed class AsObserver : ISubscriber<T>, IDisposable
{
readonly IObserver<T> observer;
ISubscription s;
bool done;
internal AsObserver(IObserver<T> observer)
{
this.observer = observer;
}
public void OnSubscribe(ISubscription s)
{
if (SubscriptionHelper.SetOnce(ref this.s, s))
{
s.Request(long.MaxValue);
}
}
public void OnNext(T t)
{
if (done)
{
return;
}
try
{
observer.OnNext(t);
}
catch (Exception ex)
{
ExceptionHelper.ThrowIfFatal(ex);
s.Cancel();
OnError(ex);
}
}
public void OnError(Exception e)
{
if (done)
{
ExceptionHelper.OnErrorDropped(e);
return;
}
done = true;
try
{
observer.OnError(e);
}
catch (Exception ex)
{
ExceptionHelper.ThrowIfFatal(ex);
ExceptionHelper.OnErrorDropped(new AggregateException(e, ex));
}
}
public void OnComplete()
{
if (done)
{
return;
}
done = true;
try
{
observer.OnCompleted();
}
catch (Exception ex)
{
ExceptionHelper.ThrowOrDrop(ex);
}
}
public void Dispose()
{
SubscriptionHelper.Cancel(ref s);
}
}
}
}