【问题标题】:Dart: How to test if stream emits elements at a certain time?Dart:如何测试流是否在特定时间发出元素?
【发布时间】:2015-11-27 12:44:18
【问题描述】:

我尝试测试一个函数Stream transform(Stream input)。如何测试返回的流是否在特定时间发出元素?

RxJS (JavaScript) 中,我可以使用TestScheduler 在某个时间在输入流上发出元素,并测试它们是否在某个时间在输出流上发出。在这个example中,transform函数被传递给scheduler.startWithCreate

var scheduler = new Rx.TestScheduler();

// Create hot observable which will start firing
var xs = scheduler.createHotObservable(
  onNext(150, 1),
  onNext(210, 2),
  onNext(220, 3),
  onCompleted(230)
);

// Note we'll start at 200 for subscribe, hence missing the 150 mark
var res = scheduler.startWithCreate(function () {
  return xs.map(function (x) { return x * x });
});

// Implement collection assertion
collectionAssert.assertEqual(res.messages, [
  onNext(210, 4),
  onNext(220, 9),
  onCompleted(230)
]);

// Check for subscribe/unsubscribe
collectionAssert.assertEqual(xs.subscriptions, [
  subscribe(200, 230)
]);

【问题讨论】:

  • 你需要这个做什么?是150210220、...、毫秒吗?多少偏差是可以接受的?检查事件以特定顺序到达是一个明显的要求,但在异步处理中没有任何时间保证。我没有仔细查看pub.dartlang.org/packages/scheduled_test。也许它包含您正在寻找的东西。
  • 我想测试流值是否正确延迟。偏差取决于延迟。 150, 210, 220, ...,在示例中可以假定为毫秒。实际上,TestScheduler 模拟了时间。因此,这次测试不会持续。因此,您甚至可以在一秒钟内测试一个小时的延迟。
  • 你不能像在 JS 中那样在 Dart 中伪造时间,因为你不能覆盖 DateTime.nowStopwatch 构造函数。如果您的程序愿意,您的程序总能获得正确的时间。您可以使用Zone 拦截计时器,但这只会对您有很大帮助。至于检查流,测试它们何时输出的方法是监听流并检查何时收到事件,我认为它必须比这更花哨。我始终同意 Günther Zöchbauer 的观点,即期望准确的时间是不安全的,您可以被任何其他执行代码任意延迟。
  • 我担心这一点。在这种情况下,10 毫秒的偏差应该是可以接受的。如何测试流元素是否在特定时间窗口内发出?

标签: unit-testing testing stream dart


【解决方案1】:

更新:将我的代码作为一个名为 stream_test_scheduler 的包发布。

此代码的工作方式类似于 RxJS 中的 TestScheduler,但它使用实时(毫秒)而不是虚拟时间,因为您无法在 Dart 中伪造时间(请参阅 Irn's comment)。您可以将最大偏差传递给匹配器。在这个例子中我使用了 20 毫秒。但偏差各不相同。您可能必须在另一个测试或另一个(更快/更慢)系统上使用不同的最大偏差值。

编辑:我将示例更改为延迟变换函数,它是 stream_ext 包的 delay 函数的较短(较少可配置/参数)版本。测试检查元素是否延迟一秒。

import 'dart:async';

import 'package:test/test.dart';

// To test:
/// Modified version of <https://github.com/theburningmonk/stream_ext/wiki/delay>
Stream delay(Stream input, Duration duration) {
  var controller = new StreamController.broadcast(sync : true);
  delayCall(Function f, [Iterable args]) => args == null
      ? new Timer(duration, f)
      : new Timer(duration, () => Function.apply(f, args));
  input.listen(
      (x) => delayCall(_tryAdd, [controller, x]),
      onError : (ex) => delayCall(_tryAddError, [ex]),
      onDone  : () => delayCall(_tryClose, [controller])
  );
  return controller.stream;
}

_tryAdd(StreamController controller, event) {
  if (!controller.isClosed) controller.add(event);
}

_tryAddError(StreamController controller, err) {
  if (!controller.isClosed) controller.addError(err);
}

_tryClose(StreamController controller) {
  if (!controller.isClosed) controller.close();
}

main() async {
  test('delay preserves relative time intervals between the values', () async {
    var scheduler = new TestScheduler();

    var source = scheduler.createStream([
      onNext(150, 1),
      onNext(210, 2),
      onNext(220, 3),
      onCompleted(230)
    ]);

    var result = await scheduler.startWithCreate(() => delay(source, ms(1000)));

    expect(result, equalsRecords([
      onNext(1150, 1),
      onNext(1210, 2),
      onNext(1220, 3),
      onCompleted(1230)
    ], maxDeviation: 20));
  });
}

equalsRecords(List<Record> records, {int maxDeviation: 0}) {
  return pairwiseCompare(records, (Record r1, Record r2) {
    var deviation = (r1.ticks.inMilliseconds - r2.ticks.inMilliseconds).abs();
    if (deviation > maxDeviation) {
      return false;
    }
    if (r1 is OnNextRecord && r2 is OnNextRecord) {
      return r1.value == r2.value;
    }
    if (r1 is OnErrorRecord && r2 is OnErrorRecord) {
      return r1.exception == r2.exception;
    }
    return (r1 is OnCompletedRecord && r2 is OnCompletedRecord);
  }, 'equal with deviation of ${maxDeviation}ms to');
}

class TestScheduler {
  final SchedulerTasks _tasks;

  TestScheduler() : _tasks = new SchedulerTasks();

  Stream createStream(List<Record> records) {
    final controller = new StreamController(sync: true);
    _tasks.add(controller, records);
    return controller.stream;
  }

  Future<List<Record>> startWithCreate(Stream createStream()) {
    final completer = new Completer<List<Record>>();
    final records = <Record>[];
    final start = new DateTime.now();
    int timeStamp() {
      final current = new DateTime.now();
      return current.difference(start).inMilliseconds;
    }
    createStream().listen(
      (event) => records.add(onNext(timeStamp(), event)),
      onError: (exception) => records.add(onError(timeStamp(), exception)),
      onDone: () {
        records.add(onCompleted(timeStamp()));
        completer.complete(records);
      }
    );
    _tasks.run();
    return completer.future;
  }
}

class SchedulerTasks {
  Map<Record, StreamController> _controllers = {};
  List<Record> _records = [];

  void add(StreamController controller, List<Record> records) {
    for (var record in records) {
      _controllers[record] = controller;
    }
    _records.addAll(records);
  }

  void run() {
    _records.sort();
    for (var record in _records) {
      final controller = _controllers[record];
      new Future.delayed(record.ticks, () {
        if (record is OnNextRecord) {
          controller.add(record.value);
        } else if (record is OnErrorRecord) {
          controller.addError(record.exception);
        } else if (record is OnCompletedRecord) {
          controller.close();
        }
      });
    }
  }
}

onNext(int ticks, int value) => new OnNextRecord(ms(ticks), value);

onCompleted(int ticks) => new OnCompletedRecord(ms(ticks));

onError(int ticks, exception) => new OnErrorRecord(ms(ticks), exception);

Duration ms(int milliseconds) => new Duration(milliseconds: milliseconds);

abstract class Record implements Comparable {
  final Duration ticks;
  Record(this.ticks);

  @override
  int compareTo(other) => ticks.compareTo(other.ticks);
}

class OnNextRecord extends Record {
  final value;
  OnNextRecord(Duration ticks, this.value) : super (ticks);

  @override
  String toString() => 'onNext($value)@${ticks.inMilliseconds}';
}

class OnErrorRecord extends Record {
  final exception;
  OnErrorRecord(Duration ticks, this.exception) : super (ticks);

  @override
  String toString() => 'onError($exception)@${ticks.inMilliseconds}';
}

class OnCompletedRecord extends Record {
  OnCompletedRecord(Duration ticks) : super (ticks);

  @override
  String toString() => 'onCompleted()@${ticks.inMilliseconds}';
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-04
    • 2019-10-15
    • 1970-01-01
    • 2019-09-15
    • 2015-11-30
    • 1970-01-01
    相关资源
    最近更新 更多