123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367 |
-
- #if (NET_4_6 || NET_STANDARD_2_0)
- using System;
- using System.Threading.Tasks;
- using System.Threading;
- namespace UniRx
- {
-
-
-
- public static class TaskObservableExtensions
- {
-
-
-
-
-
-
-
- public static IObservable<Unit> ToObservable(this Task task)
- {
- if (task == null)
- throw new ArgumentNullException("task");
- return ToObservableImpl(task, null);
- }
-
-
-
-
-
-
-
-
- public static IObservable<Unit> ToObservable(this Task task, IScheduler scheduler)
- {
- if (task == null)
- throw new ArgumentNullException("task");
- if (scheduler == null)
- throw new ArgumentNullException("scheduler");
- return ToObservableImpl(task, scheduler);
- }
- private static IObservable<Unit> ToObservableImpl(Task task, IScheduler scheduler)
- {
- var res = default(IObservable<Unit>);
- if (task.IsCompleted)
- {
- scheduler = scheduler ?? Scheduler.Immediate;
- switch (task.Status)
- {
- case TaskStatus.RanToCompletion:
- res = Observable.Return<Unit>(Unit.Default, scheduler);
- break;
- case TaskStatus.Faulted:
- res = Observable.Throw<Unit>(task.Exception.InnerException, scheduler);
- break;
- case TaskStatus.Canceled:
- res = Observable.Throw<Unit>(new TaskCanceledException(task), scheduler);
- break;
- }
- }
- else
- {
-
-
-
- res = ToObservableSlow(task, scheduler);
- }
- return res;
- }
- private static IObservable<Unit> ToObservableSlow(Task task, IScheduler scheduler)
- {
- var subject = new AsyncSubject<Unit>();
- var options = GetTaskContinuationOptions(scheduler);
- task.ContinueWith(t => ToObservableDone(task, subject), options);
- return ToObservableResult(subject, scheduler);
- }
- private static void ToObservableDone(Task task, IObserver<Unit> subject)
- {
- switch (task.Status)
- {
- case TaskStatus.RanToCompletion:
- subject.OnNext(Unit.Default);
- subject.OnCompleted();
- break;
- case TaskStatus.Faulted:
- subject.OnError(task.Exception.InnerException);
- break;
- case TaskStatus.Canceled:
- subject.OnError(new TaskCanceledException(task));
- break;
- }
- }
-
-
-
-
-
-
-
-
- public static IObservable<TResult> ToObservable<TResult>(this Task<TResult> task)
- {
- if (task == null)
- throw new ArgumentNullException("task");
- return ToObservableImpl(task, null);
- }
-
-
-
-
-
-
-
-
-
- public static IObservable<TResult> ToObservable<TResult>(this Task<TResult> task, IScheduler scheduler)
- {
- if (task == null)
- throw new ArgumentNullException("task");
- if (scheduler == null)
- throw new ArgumentNullException("scheduler");
- return ToObservableImpl(task, scheduler);
- }
- private static IObservable<TResult> ToObservableImpl<TResult>(Task<TResult> task, IScheduler scheduler)
- {
- var res = default(IObservable<TResult>);
- if (task.IsCompleted)
- {
- scheduler = scheduler ?? Scheduler.Immediate;
- switch (task.Status)
- {
- case TaskStatus.RanToCompletion:
- res = Observable.Return<TResult>(task.Result, scheduler);
- break;
- case TaskStatus.Faulted:
- res = Observable.Throw<TResult>(task.Exception.InnerException, scheduler);
- break;
- case TaskStatus.Canceled:
- res = Observable.Throw<TResult>(new TaskCanceledException(task), scheduler);
- break;
- }
- }
- else
- {
-
-
-
- res = ToObservableSlow(task, scheduler);
- }
- return res;
- }
- private static IObservable<TResult> ToObservableSlow<TResult>(Task<TResult> task, IScheduler scheduler)
- {
- var subject = new AsyncSubject<TResult>();
- var options = GetTaskContinuationOptions(scheduler);
- task.ContinueWith(t => ToObservableDone(task, subject), options);
- return ToObservableResult(subject, scheduler);
- }
- private static void ToObservableDone<TResult>(Task<TResult> task, IObserver<TResult> subject)
- {
- switch (task.Status)
- {
- case TaskStatus.RanToCompletion:
- subject.OnNext(task.Result);
- subject.OnCompleted();
- break;
- case TaskStatus.Faulted:
- subject.OnError(task.Exception.InnerException);
- break;
- case TaskStatus.Canceled:
- subject.OnError(new TaskCanceledException(task));
- break;
- }
- }
- private static TaskContinuationOptions GetTaskContinuationOptions(IScheduler scheduler)
- {
- var options = TaskContinuationOptions.None;
- if (scheduler != null)
- {
-
-
-
-
-
-
-
-
-
-
- options |= TaskContinuationOptions.ExecuteSynchronously;
- }
- return options;
- }
- private static IObservable<TResult> ToObservableResult<TResult>(AsyncSubject<TResult> subject, IScheduler scheduler)
- {
- if (scheduler != null)
- {
- return subject.ObserveOn(scheduler);
- }
- else
- {
- return subject.AsObservable();
- }
- }
-
-
-
-
-
-
-
- public static Task<TResult> ToTask<TResult>(this IObservable<TResult> observable)
- {
- if (observable == null)
- throw new ArgumentNullException("observable");
- return observable.ToTask(new CancellationToken(), null);
- }
-
-
-
-
-
-
-
-
- public static Task<TResult> ToTask<TResult>(this IObservable<TResult> observable, object state)
- {
- if (observable == null)
- throw new ArgumentNullException("observable");
- return observable.ToTask(new CancellationToken(), state);
- }
-
-
-
-
-
-
-
-
- public static Task<TResult> ToTask<TResult>(this IObservable<TResult> observable, CancellationToken cancellationToken)
- {
- if (observable == null)
- throw new ArgumentNullException("observable");
- return observable.ToTask(cancellationToken, null);
- }
-
-
-
-
-
-
-
-
-
- public static Task<TResult> ToTask<TResult>(this IObservable<TResult> observable, CancellationToken cancellationToken, object state)
- {
- if (observable == null)
- throw new ArgumentNullException("observable");
- var hasValue = false;
- var lastValue = default(TResult);
- var tcs = new TaskCompletionSource<TResult>(state);
- var disposable = new SingleAssignmentDisposable();
- var ctr = default(CancellationTokenRegistration);
- if (cancellationToken.CanBeCanceled)
- {
- ctr = cancellationToken.Register(() =>
- {
- disposable.Dispose();
- tcs.TrySetCanceled(cancellationToken);
- });
- }
- var taskCompletionObserver = Observer.Create<TResult>(
- value =>
- {
- hasValue = true;
- lastValue = value;
- },
- ex =>
- {
- tcs.TrySetException(ex);
- ctr.Dispose();
- disposable.Dispose();
- },
- () =>
- {
- if (hasValue)
- tcs.TrySetResult(lastValue);
- else
- tcs.TrySetException(new InvalidOperationException("Strings_Linq.NO_ELEMENTS"));
- ctr.Dispose();
- disposable.Dispose();
- }
- );
-
-
-
-
-
- try
- {
-
-
-
-
-
-
-
- disposable.Disposable = observable.Subscribe(taskCompletionObserver);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- }
- return tcs.Task;
- }
- }
- }
- #endif
|