Timestamp.cs 1.5 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950
  1. using System;
  2. namespace UniRx.Operators
  3. {
  4. internal class TimestampObservable<T> : OperatorObservableBase<Timestamped<T>>
  5. {
  6. readonly IObservable<T> source;
  7. readonly IScheduler scheduler;
  8. public TimestampObservable(IObservable<T> source, IScheduler scheduler)
  9. : base(scheduler == Scheduler.CurrentThread || source.IsRequiredSubscribeOnCurrentThread())
  10. {
  11. this.source = source;
  12. this.scheduler = scheduler;
  13. }
  14. protected override IDisposable SubscribeCore(IObserver<Timestamped<T>> observer, IDisposable cancel)
  15. {
  16. return source.Subscribe(new Timestamp(this, observer, cancel));
  17. }
  18. class Timestamp : OperatorObserverBase<T, Timestamped<T>>
  19. {
  20. readonly TimestampObservable<T> parent;
  21. public Timestamp(TimestampObservable<T> parent, IObserver<Timestamped<T>> observer, IDisposable cancel)
  22. : base(observer, cancel)
  23. {
  24. this.parent = parent;
  25. }
  26. public override void OnNext(T value)
  27. {
  28. base.observer.OnNext(new Timestamped<T>(value, parent.scheduler.Now));
  29. }
  30. public override void OnError(Exception error)
  31. {
  32. try { observer.OnError(error); }
  33. finally { Dispose(); }
  34. }
  35. public override void OnCompleted()
  36. {
  37. try { observer.OnCompleted(); }
  38. finally { Dispose(); }
  39. }
  40. }
  41. }
  42. }