| Index: pkg/barback/lib/src/phase.dart
|
| diff --git a/pkg/barback/lib/src/phase.dart b/pkg/barback/lib/src/phase.dart
|
| index 17d32d8e2c2fdc110572fb9deaa1226ea93b1d87..b6aafc9bade58459e17ac49ce3beb6f61684085a 100644
|
| --- a/pkg/barback/lib/src/phase.dart
|
| +++ b/pkg/barback/lib/src/phase.dart
|
| @@ -45,7 +45,7 @@ class Phase {
|
| /// The transformers that can access [inputs].
|
| ///
|
| /// Their outputs will be available to the next phase.
|
| - final Set<Transformer> _transformers;
|
| + final _transformers = new Set<Transformer>();
|
|
|
| /// The groups for this phase.
|
| final _groups = new Map<TransformerGroup, GroupRunner>();
|
| @@ -77,21 +77,22 @@ class Phase {
|
| /// so, it's been forwarded unmodified.
|
| final _inputOrigins = new Multiset<AssetNode>();
|
|
|
| - /// A stream that emits an event whenever this phase becomes dirty and needs
|
| - /// to be run.
|
| + /// A stream that emits an event whenever [this] is no longer dirty.
|
| ///
|
| - /// This may emit events when the phase was already dirty or while processing
|
| - /// transforms. Events are emitted synchronously to ensure that the dirty
|
| - /// state is thoroughly propagated as soon as any assets are changed.
|
| - Stream get onDirty => _onDirtyPool.stream;
|
| - final _onDirtyPool = new StreamPool.broadcast();
|
| + /// This is synchronous in order to guarantee that it will emit an event as
|
| + /// soon as [isDirty] flips from `true` to `false`.
|
| + Stream get onDone => _onDoneController.stream;
|
| + final _onDoneController = new StreamController.broadcast(sync: true);
|
|
|
| - /// A controller whose stream feeds into [_onDirtyPool].
|
| + /// A stream that emits any new assets emitted by [this].
|
| ///
|
| - /// This is used whenever an input is added or transforms are changed.
|
| - final _onDirtyController = new StreamController.broadcast(sync: true);
|
| + /// Assets are emitted synchronously to ensure that any changes are thoroughly
|
| + /// propagated as soon as they occur. Only a phase with no [next] phase will
|
| + /// emit assets.
|
| + Stream<AssetNode> get onAsset => _onAssetController.stream;
|
| + final _onAssetController = new StreamController<AssetNode>(sync: true);
|
|
|
| - /// Whether this phase is dirty and needs to be run.
|
| + /// Whether [this] is dirty and still has more processing to do.
|
| bool get isDirty => _inputs.values.any((input) => input.isDirty) ||
|
| _groups.values.any((group) => group.isDirty);
|
|
|
| @@ -117,20 +118,10 @@ class Phase {
|
| // TODO(nweiz): Rather than passing the cascade and the phase everywhere,
|
| // create an interface that just exposes [getInput]. Emit errors via
|
| // [AssetNode]s.
|
| - Phase(AssetCascade cascade, Iterable transformers, String location)
|
| - : this._(cascade, transformers, location, 0);
|
| + Phase(AssetCascade cascade, String location)
|
| + : this._(cascade, location, 0);
|
|
|
| - Phase._(this.cascade, Iterable transformers, this._location, this._index)
|
| - : _transformers = transformers.where((op) => op is Transformer).toSet() {
|
| - _onDirtyPool.add(_onDirtyController.stream);
|
| -
|
| - for (var group in transformers.where((op) => op is TransformerGroup)) {
|
| - var runner = new GroupRunner(cascade, group, "$_location.$_index");
|
| - _groups[group] = runner;
|
| - _onDirtyPool.add(runner.onDirty);
|
| - _onLogPool.add(runner.onLog);
|
| - }
|
| - }
|
| + Phase._(this.cascade, this._location, this._index);
|
|
|
| /// Adds a new asset as an input for this phase.
|
| ///
|
| @@ -151,12 +142,7 @@ class Phase {
|
| // there's one additional channel for the non-grouped transformers.
|
| var forwarder = new PhaseForwarder(_groups.length + 1);
|
| _forwarders[node.id] = forwarder;
|
| - forwarder.onForwarding.listen((asset) {
|
| - _addOutput(asset);
|
| -
|
| - var exception = _outputs[asset.id].collisionException;
|
| - if (exception != null) cascade.reportError(exception);
|
| - });
|
| + forwarder.onAsset.listen(_handleOutputWithoutForwarder);
|
|
|
| _inputOrigins.add(node.origin);
|
| var input = new PhaseInput(this, node, _transformers, "$_location.$_index");
|
| @@ -165,10 +151,13 @@ class Phase {
|
| _inputOrigins.remove(node.origin);
|
| _inputs.remove(node.id);
|
| _forwarders.remove(node.id).remove();
|
| + if (!isDirty) _onDoneController.add(null);
|
| });
|
| - _onDirtyPool.add(input.onDirty);
|
| - _onDirtyController.add(null);
|
| + input.onAsset.listen(_handleOutput);
|
| _onLogPool.add(input.onLog);
|
| + input.onDone.listen((_) {
|
| + if (!isDirty) _onDoneController.add(null);
|
| + });
|
|
|
| for (var group in _groups.values) {
|
| group.addInput(node);
|
| @@ -201,8 +190,6 @@ class Phase {
|
|
|
| /// Set this phase's transformers to [transformers].
|
| void updateTransformers(Iterable transformers) {
|
| - _onDirtyController.add(null);
|
| -
|
| var actualTransformers = transformers.where((op) => op is Transformer);
|
| _transformers.clear();
|
| _transformers.addAll(actualTransformers);
|
| @@ -220,8 +207,11 @@ class Phase {
|
| for (var added in newGroups.difference(oldGroups)) {
|
| var runner = new GroupRunner(cascade, added, "$_location.$_index");
|
| _groups[added] = runner;
|
| - _onDirtyPool.add(runner.onDirty);
|
| + runner.onAsset.listen(_handleOutput);
|
| _onLogPool.add(runner.onLog);
|
| + runner.onDone.listen((_) {
|
| + if (!isDirty) _onDoneController.add(null);
|
| + });
|
| for (var input in _inputs.values) {
|
| runner.addInput(input.input);
|
| }
|
| @@ -244,12 +234,12 @@ class Phase {
|
| }
|
| }
|
|
|
| - /// Add a new phase after this one with [transformers].
|
| + /// Add a new phase after this one.
|
| ///
|
| /// This may only be called on a phase with no phase following it.
|
| - Phase addPhase(Iterable transformers) {
|
| + Phase addPhase() {
|
| assert(_next == null);
|
| - _next = new Phase._(cascade, transformers, _location, _index + 1);
|
| + _next = new Phase._(cascade, _location, _index + 1);
|
| for (var output in _outputs.values.toList()) {
|
| // Remove [output]'s listeners because now they should get the asset from
|
| // [_next], rather than this phase. Any transforms consuming [output] will
|
| @@ -270,7 +260,7 @@ class Phase {
|
| for (var group in _groups.values) {
|
| group.remove();
|
| }
|
| - _onDirtyPool.close();
|
| + _onAssetController.close();
|
| _onLogPool.close();
|
| }
|
|
|
| @@ -281,60 +271,40 @@ class Phase {
|
| _next = null;
|
| }
|
|
|
| - /// Processes this phase.
|
| - ///
|
| - /// Returns a future that completes when processing is done. If there is
|
| - /// nothing to process, returns `null`.
|
| - Future process() {
|
| - if (!isDirty) return null;
|
| -
|
| - var outputIds = new Set<AssetId>();
|
| - void _handleOutputs(Set<AssetNode> outputs) {
|
| - for (var asset in outputs) {
|
| - if (_inputOrigins.contains(asset.origin)) {
|
| - _forwarders[asset.id].addIntermediateAsset(asset);
|
| - continue;
|
| - }
|
| -
|
| - outputIds.add(asset.id);
|
| - _addOutput(asset);
|
| - }
|
| + /// Add [asset] as an output of this phase.
|
| + void _handleOutput(AssetNode asset) {
|
| + if (_inputOrigins.contains(asset.origin)) {
|
| + _forwarders[asset.id].addIntermediateAsset(asset);
|
| + } else {
|
| + _handleOutputWithoutForwarder(asset);
|
| }
|
| -
|
| - var outputFutures = [];
|
| - outputFutures.addAll(_inputs.values.map((input) {
|
| - if (!input.isDirty) return new Future.value(new Set());
|
| - return input.process().then(_handleOutputs);
|
| - }));
|
| - outputFutures.addAll(_groups.values.map((group) {
|
| - if (!group.isDirty) return new Future.value(new Set());
|
| - return group.process().then(_handleOutputs);
|
| - }));
|
| -
|
| - return Future.wait(outputFutures).then((_) {
|
| - // Report collisions in a deterministic order.
|
| - outputIds = outputIds.toList();
|
| - outputIds.sort((a, b) => a.compareTo(b));
|
| - for (var id in outputIds) {
|
| - // It's possible the output was removed before other transforms in this
|
| - // phase finished.
|
| - if (!_outputs.containsKey(id)) continue;
|
| - var exception = _outputs[id].collisionException;
|
| - if (exception != null) cascade.reportError(exception);
|
| - }
|
| - });
|
| }
|
|
|
| - /// Add [asset] as an output of this phase.
|
| - void _addOutput(AssetNode asset) {
|
| + /// Add [asset] as an output of this phase without checking if it's a
|
| + /// forwarded asset.
|
| + void _handleOutputWithoutForwarder(AssetNode asset) {
|
| if (_outputs.containsKey(asset.id)) {
|
| _outputs[asset.id].add(asset);
|
| } else {
|
| _outputs[asset.id] = new PhaseOutput(this, asset, "$_location.$_index");
|
| - _outputs[asset.id].onAsset.listen((output) {
|
| - if (_next != null) _next.addInput(output);
|
| - }, onDone: () => _outputs.remove(asset.id));
|
| - if (_next != null) _next.addInput(_outputs[asset.id].output);
|
| + _outputs[asset.id].onAsset.listen(_emit,
|
| + onDone: () => _outputs.remove(asset.id));
|
| + _emit(_outputs[asset.id].output);
|
| + }
|
| +
|
| + var exception = _outputs[asset.id].collisionException;
|
| + if (exception != null) cascade.reportError(exception);
|
| + }
|
| +
|
| + /// Emit [asset] as an output of this phase.
|
| + ///
|
| + /// This should be called after [_handleOutput], so that collisions are
|
| + /// resolved.
|
| + void _emit(AssetNode asset) {
|
| + if (_next != null) {
|
| + _next.addInput(asset);
|
| + } else {
|
| + _onAssetController.add(asset);
|
| }
|
| }
|
|
|
|
|