Chromium Code Reviews| 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..094e54983cb9b93505580d95f78511e59912560f 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(isDirty); |
|
Bob Nystrom
2014/03/05 22:13:25
.add(null);
nweiz
2014/03/06 00:29:08
Done.
|
| + }); |
| 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); |
| } |
| } |