单元测试失败,Observable.FromAsync和Observable.Switch



我在测试一个使用Observable.FromAsync<T>()Observable.Switch<T>()的类时遇到问题。它所做的是等待一个可观察到的触发器产生一个值,然后启动一个异步操作,最后在一个输出序列中回收所有操作的结果。它的要点是:

var outputStream = triggerStream
  .Select(_ => Observable
    .FromAsync(token => taskProducer.DoSomethingAsync(token)))
  .Switch();

我用最基本的部分进行了一些健全性检查测试,以了解发生了什么,以下是测试结果,并在评论中给出:

class test_with_rx : nspec
{
  void Given_async_task_and_switch()
  {
    Subject<Unit> triggerStream = null;
    TaskCompletionSource<long> taskDriver = null;
    ITestableObserver<long> testObserver = null;
    IDisposable subscription = null;
    before = () =>
    {
      TestScheduler scheduler = new TestScheduler();
      testObserver = scheduler.CreateObserver<long>();
      triggerStream = new Subject<Unit>();
      taskDriver = new TaskCompletionSource<long>();
      // build stream under test
      IObservable<long> streamUnderTest = triggerStream
        .Select(_ => Observable
          .FromAsync(token => taskDriver.Task))
        .Switch();
      /* Also tried with this Switch() overload
      IObservable<long> streamUnderTest = triggerStream
          .Select(_ => taskDriver.Task)
          .Switch(); */
      subscription = streamUnderTest.Subscribe(testObserver);
    };
    context["Before trigger"] = () =>
    {
      it["Should not notify"] = () => testObserver.Messages.Count.Should().Be(0);
      // PASSED
    };
    context["After trigger"] = () =>
    {
      before = () => triggerStream.OnNext(Unit.Default);
      context["When task completes"] = () =>
      {
        long result = -1;
        before = () =>
        {
          taskDriver.SetResult(result);
          //taskDriver.Task.Wait();  // tried with this too
        };
        it["Should notify once"] = () => testObserver.Messages.Count.Should().Be(1);
        // FAILED: expected 1, actual 0
        it["Should notify task result"] = () => testObserver.Messages[0].Value.Value.Should().Be(result);
        // FAILED: of course, index out of bound
      };
    };
    after = () =>
    {
      taskDriver.TrySetCanceled();
      taskDriver.Task.Dispose();
      subscription.Dispose();
    };
  }
}

在我对mock所做的其他测试中,我可以看到传递给FromAsync的Func实际上被调用了(例如taskProducer.DoSomethingAsync(token)(,但之后看起来什么都没有了,输出流也没有产生值。

我还试着在达到预期之前插入一些Task.Delay(x).Wait()taskDriver.Task.Wait(),但没有成功。

我读过这个SO线程,我知道调度器,但乍一看,我觉得我不需要它们,没有使用ObserveOn()。我错了吗?我错过了什么?TA

为了完整性,测试框架是NSpec,断言库是FluentAssessments。

您遇到的是一起测试Rx和TPL的情况。这里可以找到详尽的解释,但我会尝试为您的特定代码提供建议。

基本上,您的代码运行良好,但您的测试不正常。Observable.FromAsync将在所提供的任务上转换为ContinueWith,该任务将在任务池上执行,因此是异步的。

修复测试的多种方法:(从丑陋到复杂(

  1. 结果设置后睡眠(注意等待不起作用,因为等待不等待继续(

    taskDriver.SetResult(result);
    Thread.Sleep(50);
    
  2. 在执行FromAsync之前设置结果(因为如果任务完成,FromAsync将立即返回IObservable,也就是跳过ContinueWith(

    taskDriver.SetResult(result);
    triggerStream.OnNext(Unit.Default);
    
  3. 用可测试的替代方案替换FromAsync,例如

    public static IObservable<T> ToObservable<T>(Task<T> task, TaskScheduler scheduler)
    {
        if (task.IsCompleted)
        {
            return task.ToObservable();
        }
        else
        {
            AsyncSubject<T> asyncSubject = new AsyncSubject<T>();
            task.ContinueWith(t => task.ToObservable().Subscribe(asyncSubject), scheduler);
            return asyncSubject.AsObservable<T>();
        }
    }
    

(使用同步TaskScheduler或可测试的TaskScheduler(

相关内容

  • 没有找到相关文章

最新更新