Permalink
Browse files

Merge pull request #162 from garuma/tpl-dataflow-blocks

@marek - TPL Dataflow part 4
  • Loading branch information...
marek-safar committed Aug 23, 2011
2 parents cc8c07f + 639c8a4 commit c4959f4f3b79218735188088584ba5e54df3d97b
Showing with 3,244 additions and 1 deletion.
  1. +15 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow.dll.sources
  2. +111 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/ActionBlock.cs
  3. +171 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/BatchBlock.cs
  4. +138 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/BroadcastBlock.cs
  5. +143 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/BufferBlock.cs
  6. +113 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/ChooserBlock.cs
  7. +221 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/DataflowBlock.cs
  8. +211 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/JoinBlock.cs
  9. +233 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/JoinBlock`3.cs
  10. +2 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/MessageOutgoingQueue.cs
  11. +94 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/ObservableDataflowBlock.cs
  12. +58 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/ObserverDataflowBlock.cs
  13. +94 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/PropagatorWrapperBlock.cs
  14. +116 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/ReceiveBlock.cs
  15. +156 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/TransformBlock.cs
  16. +156 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/TransformManyBlock.cs
  17. +168 −0 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow/WriteOnceBlock.cs
  18. +11 −1 mcs/class/System.Threading.Tasks.Dataflow/System.Threading.Tasks.Dataflow_test.dll.sources
  19. +70 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/ActionBlockTest.cs
  20. +129 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/BatchBlockTest.cs
  21. +87 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/BroadcastBlockTest.cs
  22. +97 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/BufferBlockTest.cs
  23. +89 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/CompletionTest.cs
  24. +143 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/DataflowBlockTest.cs
  25. +63 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/JoinBlockTest.cs
  26. +69 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/JoinBlock`3Test.cs
  27. +75 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/TransformBlockTest.cs
  28. +78 −0 ...ss/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/TransformManyBlockTest.cs
  29. +133 −0 mcs/class/System.Threading.Tasks.Dataflow/Test/System.Threading.Tasks.Dataflow/WriteOnceBlockTest.cs
@@ -21,3 +21,18 @@ System.Threading.Tasks.Dataflow/MessageVault.cs
System.Threading.Tasks.Dataflow/PassingMessageBox.cs
System.Threading.Tasks.Dataflow/TargetBuffer.cs
../corlib/System.Threading/AtomicBoolean.cs
+System.Threading.Tasks.Dataflow/ActionBlock.cs
+System.Threading.Tasks.Dataflow/BatchBlock.cs
+System.Threading.Tasks.Dataflow/BroadcastBlock.cs
+System.Threading.Tasks.Dataflow/BufferBlock.cs
+System.Threading.Tasks.Dataflow/ChooserBlock.cs
+System.Threading.Tasks.Dataflow/DataflowBlock.cs
+System.Threading.Tasks.Dataflow/JoinBlock.cs
+System.Threading.Tasks.Dataflow/JoinBlock`3.cs
+System.Threading.Tasks.Dataflow/ObservableDataflowBlock.cs
+System.Threading.Tasks.Dataflow/ObserverDataflowBlock.cs
+System.Threading.Tasks.Dataflow/PropagatorWrapperBlock.cs
+System.Threading.Tasks.Dataflow/ReceiveBlock.cs
+System.Threading.Tasks.Dataflow/TransformBlock.cs
+System.Threading.Tasks.Dataflow/TransformManyBlock.cs
+System.Threading.Tasks.Dataflow/WriteOnceBlock.cs
@@ -0,0 +1,111 @@
+// ActionBlock.cs
+//
+// Copyright (c) 2011 Jérémie "garuma" Laval
+//
+// Permission is hereby granted, free of charge, to any person obtaining a copy
+// of this software and associated documentation files (the "Software"), to deal
+// in the Software without restriction, including without limitation the rights
+// to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
+// copies of the Software, and to permit persons to whom the Software is
+// furnished to do so, subject to the following conditions:
+//
+// The above copyright notice and this permission notice shall be included in
+// all copies or substantial portions of the Software.
+//
+// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
+// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
+// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
+// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
+// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
+// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
+// THE SOFTWARE.
+//
+//
+
+
+using System;
+using System.Threading;
+using System.Threading.Tasks;
+using System.Collections.Concurrent;
+
+namespace System.Threading.Tasks.Dataflow
+{
+ public sealed class ActionBlock<TInput> : ITargetBlock<TInput>, IDataflowBlock
+ {
+ static readonly ExecutionDataflowBlockOptions defaultOptions = new ExecutionDataflowBlockOptions ();
+
+ CompletionHelper compHelper = CompletionHelper.GetNew ();
+ BlockingCollection<TInput> messageQueue = new BlockingCollection<TInput> ();
+ ExecutingMessageBox<TInput> messageBox;
+ Action<TInput> action;
+ ExecutionDataflowBlockOptions dataflowBlockOptions;
+
+
+ public ActionBlock (Action<TInput> action) : this (action, defaultOptions)
+ {
+
+ }
+
+ public ActionBlock (Action<TInput> action, ExecutionDataflowBlockOptions dataflowBlockOptions)
+ {
+ if (action == null)
+ throw new ArgumentNullException ("action");
+ if (dataflowBlockOptions == null)
+ throw new ArgumentNullException ("dataflowBlockOptions");
+
+ this.action = action;
+ this.dataflowBlockOptions = dataflowBlockOptions;
+ this.messageBox = new ExecutingMessageBox<TInput> (messageQueue, compHelper, () => true, ProcessQueue, dataflowBlockOptions);
+ }
+
+ [MonoTODO]
+ public ActionBlock (Func<TInput, Task> action) : this (action, defaultOptions)
+ {
+ throw new NotImplementedException ();
+ }
+
+ [MonoTODO]
+ public ActionBlock (Func<TInput, Task> action, ExecutionDataflowBlockOptions dataflowBlockOptions)
+ {
+ throw new NotImplementedException ();
+ }
+
+ public DataflowMessageStatus OfferMessage (DataflowMessageHeader messageHeader,
+ TInput messageValue,
+ ISourceBlock<TInput> source,
+ bool consumeToAccept)
+ {
+ return messageBox.OfferMessage (this, messageHeader, messageValue, source, consumeToAccept);
+ }
+
+ void ProcessQueue ()
+ {
+ TInput data;
+ while (messageQueue.TryTake (out data))
+ action (data);
+ }
+
+ public void Complete ()
+ {
+ messageBox.Complete ();
+ }
+
+ public void Fault (Exception ex)
+ {
+ compHelper.Fault (ex);
+ }
+
+ public Task Completion {
+ get {
+ return compHelper.Completion;
+ }
+ }
+
+ public int InputCount {
+ get {
+ return messageQueue.Count;
+ }
+ }
+ }
+}
+
@@ -0,0 +1,171 @@
+// BatchBlock.cs
+//
+// Copyright (c) 2011 Jérémie "garuma" Laval
+//
+// Permission is hereby granted, free of charge, to any person obtaining a copy
+// of this software and associated documentation files (the "Software"), to deal
+// in the Software without restriction, including without limitation the rights
+// to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
+// copies of the Software, and to permit persons to whom the Software is
+// furnished to do so, subject to the following conditions:
+//
+// The above copyright notice and this permission notice shall be included in
+// all copies or substantial portions of the Software.
+//
+// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
+// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
+// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
+// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
+// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
+// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
+// THE SOFTWARE.
+//
+//
+
+
+using System;
+using System.Threading.Tasks;
+using System.Collections.Generic;
+using System.Collections.Concurrent;
+
+namespace System.Threading.Tasks.Dataflow
+{
+ public sealed class BatchBlock<T> : IPropagatorBlock<T, T[]>, ITargetBlock<T>, IDataflowBlock, ISourceBlock<T[]>, IReceivableSourceBlock<T[]>
+ {
+ static readonly DataflowBlockOptions defaultOptions = new DataflowBlockOptions ();
+
+ CompletionHelper compHelper = CompletionHelper.GetNew ();
+ BlockingCollection<T> messageQueue = new BlockingCollection<T> ();
+ MessageBox<T> messageBox;
+ MessageVault<T[]> vault;
+ DataflowBlockOptions dataflowBlockOptions;
+ readonly int batchSize;
+ int batchCount;
+ MessageOutgoingQueue<T[]> outgoing;
+ TargetBuffer<T[]> targets = new TargetBuffer<T[]> ();
+ DataflowMessageHeader headers = DataflowMessageHeader.NewValid ();
+
+ public BatchBlock (int batchSize) : this (batchSize, defaultOptions)
+ {
+
+ }
+
+ public BatchBlock (int batchSize, DataflowBlockOptions dataflowBlockOptions)
+ {
+ if (dataflowBlockOptions == null)
+ throw new ArgumentNullException ("dataflowBlockOptions");
+
+ this.batchSize = batchSize;
+ this.dataflowBlockOptions = dataflowBlockOptions;
+ this.messageBox = new PassingMessageBox<T> (messageQueue, compHelper, () => outgoing.IsCompleted, BatchProcess, dataflowBlockOptions);
+ this.outgoing = new MessageOutgoingQueue<T[]> (compHelper, () => messageQueue.IsCompleted);
+ this.vault = new MessageVault<T[]> ();
+ }
+
+ public DataflowMessageStatus OfferMessage (DataflowMessageHeader messageHeader,
+ T messageValue,
+ ISourceBlock<T> source,
+ bool consumeToAccept)
+ {
+ return messageBox.OfferMessage (this, messageHeader, messageValue, source, consumeToAccept);
+ }
+
+ public IDisposable LinkTo (ITargetBlock<T[]> target, bool unlinkAfterOne)
+ {
+ var result = targets.AddTarget (target, unlinkAfterOne);
+ outgoing.ProcessForTarget (target, this, false, ref headers);
+
+ return result;
+ }
+
+ public T[] ConsumeMessage (DataflowMessageHeader messageHeader, ITargetBlock<T[]> target, out bool messageConsumed)
+ {
+ return vault.ConsumeMessage (messageHeader, target, out messageConsumed);
+ }
+
+ public void ReleaseReservation (DataflowMessageHeader messageHeader, ITargetBlock<T[]> target)
+ {
+ vault.ReleaseReservation (messageHeader, target);
+ }
+
+ public bool ReserveMessage (DataflowMessageHeader messageHeader, ITargetBlock<T[]> target)
+ {
+ return vault.ReserveMessage (messageHeader, target);
+ }
+
+ public bool TryReceive (Predicate<T[]> filter, out T[] item)
+ {
+ return TryReceive (filter, out item);
+ }
+
+ public bool TryReceiveAll (out IList<T[]> items)
+ {
+ return outgoing.TryReceiveAll (out items);
+ }
+
+ public void TriggerBatch ()
+ {
+ int earlyBatchSize;
+ do {
+ earlyBatchSize = batchCount;
+ if (earlyBatchSize == 0)
+ return;
+ } while (Interlocked.CompareExchange (ref batchCount, 0, earlyBatchSize) != earlyBatchSize);
+
+ MakeBatch (targets.Current, earlyBatchSize);
+ }
+
+ // TODO: there can be out-of-order processing of message elements if two collections
+ // are triggered and work side by side. See if it's a problem or not.
+ void BatchProcess ()
+ {
+ ITargetBlock<T[]> target = targets.Current;
+ int current = Interlocked.Increment (ref batchCount);
+
+ if (current % batchSize != 0)
+ return;
+
+ Interlocked.Add (ref batchCount, -current);
+
+ MakeBatch (target, batchSize);
+ }
+
+ void MakeBatch (ITargetBlock<T[]> target, int size)
+ {
+ T[] batch = new T[size];
+ for (int i = 0; i < size; ++i)
+ messageQueue.TryTake (out batch[i]);
+
+ if (target == null)
+ outgoing.AddData (batch);
+ else
+ target.OfferMessage (headers.Increment (), batch, this, false);
+
+ if (!outgoing.IsEmpty && targets.Current != null)
+ outgoing.ProcessForTarget (targets.Current, this, false, ref headers);
+ }
+
+ public void Complete ()
+ {
+ messageBox.Complete ();
+ }
+
+ public void Fault (Exception ex)
+ {
+ compHelper.Fault (ex);
+ }
+
+ public Task Completion {
+ get {
+ return compHelper.Completion;
+ }
+ }
+
+ public int OutputCount {
+ get {
+ return outgoing.Count;
+ }
+ }
+ }
+}
+
Oops, something went wrong.

0 comments on commit c4959f4

Please sign in to comment.