Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(4)

Unified Diff: pkg/barback/lib/src/phase.dart

Issue 25376003: Add support for transformer clusters to barback. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: More code reivew changes Created 7 years, 2 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View side-by-side diff with in-line comments
Download patch
« no previous file with comments | « pkg/barback/lib/src/group_runner.dart ('k') | pkg/barback/lib/src/phase_forwarder.dart » ('j') | no next file with comments »
Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
Index: pkg/barback/lib/src/phase.dart
diff --git a/pkg/barback/lib/src/phase.dart b/pkg/barback/lib/src/phase.dart
index 1230bea094e36970d402547ec6b386c868c33ec0..b83061bb96a1a200fd730319aef32d4c70375978 100644
--- a/pkg/barback/lib/src/phase.dart
+++ b/pkg/barback/lib/src/phase.dart
@@ -11,11 +11,14 @@ import 'asset_cascade.dart';
import 'asset_id.dart';
import 'asset_node.dart';
import 'asset_set.dart';
+import 'group_runner.dart';
import 'errors.dart';
+import 'phase_forwarder.dart';
import 'phase_input.dart';
import 'phase_output.dart';
import 'stream_pool.dart';
import 'transformer.dart';
+import 'transformer_group.dart';
import 'utils.dart';
/// One phase in the ordered series of transformations in an [AssetCascade].
@@ -39,15 +42,38 @@ class Phase {
/// Their outputs will be available to the next phase.
final Set<Transformer> _transformers;
+ /// The groups for this phase.
+ final _groups = new Map<TransformerGroup, GroupRunner>();
+
/// The inputs for this phase.
///
/// For the first phase, these will be the source assets. For all other
/// phases, they will be the outputs from the previous phase.
final _inputs = new Map<AssetId, PhaseInput>();
+ /// The forwarders for this phase.
+ final _forwarders = new Map<AssetId, PhaseForwarder>();
+
/// The outputs for this phase.
final _outputs = new Map<AssetId, PhaseOutput>();
+ // TODO(nweiz): Don't re-calculate this on the fly all the time.
+ /// The set of all [AssetNode.origin] properties of the input assets for this
+ /// phase.
+ ///
+ /// This is used to determine which assets have been passed unmodified through
+ /// [_inputs] or [_groups]. Each input asset has a PhaseInput in [_inputs]. If
+ /// that input isn't consumed by any transformers, it will be forwarded
+ /// through the PhaseInput. However, it's possible that it was consumed by a
+ /// group, and so shouldn't be forwarded through the phase as a whole.
+ ///
+ /// In order to detect whether an output has been forwarded through a group or
+ /// a PhaseInput, we must be able to distinguish it from other outputs with
+ /// the same id. To do so, we check if its origin is in [_inputOrigins]. If
+ /// so, it's been forwarded unmodified.
+ Set<AssetNode> get _inputOrigins =>
+ _inputs.values.map((input) => input.input.origin).toSet();
+
/// A stream that emits an event whenever this phase becomes dirty and needs
/// to be run.
///
@@ -62,6 +88,10 @@ class Phase {
/// This is used whenever an input is added or transforms are changed.
final _onDirtyController = new StreamController.broadcast(sync: true);
+ /// Whether this phase is dirty and needs to be run.
+ bool get isDirty => _inputs.values.any((input) => input.isDirty) ||
+ _groups.values.any((group) => group.isDirty);
+
/// The phase after this one.
///
/// Outputs from this phase will be passed to it.
@@ -69,16 +99,24 @@ class Phase {
Phase _next;
/// Returns all currently-available output assets for this phase.
- AssetSet get availableOutputs {
- return new AssetSet.from(_outputs.values
+ Set<AssetNode> get availableOutputs {
+ return _outputs.values
.map((output) => output.output)
.where((node) => node.state.isAvailable)
- .map((node) => node.asset));
+ .toSet();
}
- Phase(this.cascade, Iterable<Transformer> transformers)
- : _transformers = transformers.toSet() {
+ // TODO(nweiz): Rather than passing the cascade and the phase everywhere,
+ // create an interface that just exposes [getInput]. Emit errors via
+ // [AssetNode]s.
+ Phase(this.cascade, Iterable transformers)
+ : _transformers = transformers.where((op) => op is Transformer).toSet() {
_onDirtyPool.add(_onDirtyController.stream);
+
+ for (var group in transformers.where((op) => op is TransformerGroup)) {
+ _groups[group] = new GroupRunner(cascade, group);
+ _onDirtyPool.add(_groups[group].onDirty);
+ }
}
/// Adds a new asset as an input for this phase.
@@ -94,11 +132,29 @@ class Phase {
void addInput(AssetNode node) {
if (_inputs.containsKey(node.id)) _inputs[node.id].remove();
+ // Each group is one channel along which an asset may be forwarded. Then
+ // 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);
+ });
+
var input = new PhaseInput(this, node, _transformers);
_inputs[node.id] = input;
- input.input.whenRemoved.then((_) => _inputs.remove(node.id));
+ input.input.whenRemoved.then((_) {
+ _inputs.remove(node.id);
+ _forwarders.remove(node.id).remove();
+ });
_onDirtyPool.add(input.onDirty);
_onDirtyController.add(null);
+
+ for (var group in _groups.values) {
+ group.addInput(node);
+ }
}
/// Gets the asset node for an input [id].
@@ -124,19 +180,41 @@ class Phase {
}
/// Set this phase's transformers to [transformers].
- void updateTransformers(Iterable<Transformer> transformers) {
+ void updateTransformers(Iterable transformers) {
_onDirtyController.add(null);
+
+ var actualTransformers = transformers.where((op) => op is Transformer);
_transformers.clear();
- _transformers.addAll(transformers);
+ _transformers.addAll(actualTransformers);
for (var input in _inputs.values) {
- input.updateTransformers(_transformers);
+ input.updateTransformers(actualTransformers);
+ }
+
+ var newGroups = transformers.where((op) => op is TransformerGroup)
+ .toSet();
+ var oldGroups = _groups.keys.toSet();
+ for (var removed in oldGroups.difference(newGroups)) {
+ _groups.remove(removed).remove();
+ }
+
+ for (var added in newGroups.difference(oldGroups)) {
+ var runner = new GroupRunner(cascade, added);
+ _groups[added] = runner;
+ _onDirtyPool.add(runner.onDirty);
+ for (var input in _inputs.values) {
+ runner.addInput(input.input);
+ }
+ }
+
+ for (var forwarder in _forwarders.values) {
+ forwarder.numChannels = _groups.length + 1;
}
}
/// Add a new phase after this one with [transformers].
///
/// This may only be called on a phase with no phase following it.
- Phase addPhase(Iterable<Transformer> transformers) {
+ Phase addPhase(Iterable transformers) {
assert(_next == null);
_next = new Phase(cascade, transformers);
for (var output in _outputs.values.toList()) {
@@ -162,6 +240,9 @@ class Phase {
for (var input in _inputs.values.toList()) {
input.remove();
}
+ for (var group in _groups.values) {
+ group.remove();
+ }
_onDirtyPool.close();
}
@@ -177,26 +258,34 @@ class Phase {
/// Returns a future that completes when processing is done. If there is
/// nothing to process, returns `null`.
Future process() {
- if (!_inputs.values.any((input) => input.isDirty)) return null;
+ if (!isDirty) return null;
var outputIds = new Set<AssetId>();
- return Future.wait(_inputs.values.map((input) {
- if (!input.isDirty) return new Future.value(new Set());
- return input.process().then((outputs) {
- for (var asset in outputs) {
- outputIds.add(asset.id);
- if (_outputs.containsKey(asset.id)) {
- _outputs[asset.id].add(asset);
- } else {
- _outputs[asset.id] = new PhaseOutput(this, asset);
- _outputs[asset.id].output.whenRemoved.then((_) {
- _outputs.remove(asset.id);
- });
- if (_next != null) _next.addInput(_outputs[asset.id].output);
- }
+ void _handleOutputs(Set<AssetNode> outputs) {
+ for (var asset in outputs) {
+ if (_inputOrigins.contains(asset.origin)) {
+ _forwarders[asset.id].addIntermediateAsset(asset);
+ continue;
}
- });
- })).then((_) {
+
+ outputIds.add(asset.id);
+ _addOutput(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((input) {
+ if (!input.isDirty) return new Future.value(new Set());
+ return input.process().then(_handleOutputs);
+ }));
+
+ // TODO(nweiz): handle pass-through.
+
+ return Future.wait(outputFutures).then((_) {
// Report collisions in a deterministic order.
outputIds = outputIds.toList();
outputIds.sort((a, b) => a.compareTo(b));
@@ -209,4 +298,17 @@ class Phase {
}
});
}
+
+ /// Add [asset] as an output of this phase.
+ void _addOutput(AssetNode asset) {
+ if (_outputs.containsKey(asset.id)) {
+ _outputs[asset.id].add(asset);
+ } else {
+ _outputs[asset.id] = new PhaseOutput(this, asset);
+ _outputs[asset.id].output.whenRemoved.then((_) {
+ _outputs.remove(asset.id);
+ });
+ if (_next != null) _next.addInput(_outputs[asset.id].output);
+ }
+ }
}
« no previous file with comments | « pkg/barback/lib/src/group_runner.dart ('k') | pkg/barback/lib/src/phase_forwarder.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698