Data pipeline refactoring
This commit is contained in:
@@ -0,0 +1,894 @@
|
||||
unit Myc.Data.Pipeline.Impl;
|
||||
|
||||
interface
|
||||
|
||||
uses
|
||||
System.SysUtils,
|
||||
System.SyncObjs,
|
||||
System.Generics.Collections,
|
||||
Myc.Signals,
|
||||
Myc.Mutable,
|
||||
Myc.Core.Notifier,
|
||||
Myc.Data.Series,
|
||||
Myc.Data.Pipeline;
|
||||
|
||||
type
|
||||
// Abstract base class for data consumers.
|
||||
TMycConsumer<T> = class abstract(TContainedObject, IConsumer<T>)
|
||||
strict private
|
||||
type
|
||||
TNull = class(TInterfacedObject, IConsumer<T>)
|
||||
protected
|
||||
function Consume(const Value: T): TState;
|
||||
end;
|
||||
class var
|
||||
FNull: IConsumer<T>;
|
||||
class constructor CreateClass;
|
||||
protected
|
||||
function Consume(const Value: T): TState; virtual; abstract;
|
||||
public
|
||||
class property Null: IConsumer<T> read FNull;
|
||||
end;
|
||||
|
||||
TMycProducer<T> = class(TInterfacedObject, IProducer<T>)
|
||||
strict private
|
||||
type
|
||||
TNull = class(TInterfacedObject, IProducer<T>)
|
||||
public
|
||||
function Link(const Consumer: IConsumer<T>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
end;
|
||||
class var
|
||||
FNull: IProducer<T>;
|
||||
class constructor CreateClass;
|
||||
private
|
||||
FListeners: TMycNotifyList<IConsumer<T>>;
|
||||
public
|
||||
constructor Create;
|
||||
destructor Destroy; override;
|
||||
function Broadcast(const Value: T): TState;
|
||||
function Link(const Consumer: IConsumer<T>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
class property Null: IProducer<T> read FNull;
|
||||
end;
|
||||
|
||||
// A consumer implementation that is owned by a controller.
|
||||
TMycGenericConsumer<T> = class(TMycConsumer<T>)
|
||||
private
|
||||
FProc: TConvertFunc<T, TState>;
|
||||
protected
|
||||
function Consume(const Value: T): TState; override; final;
|
||||
public
|
||||
constructor Create(const Controller: IInterface; const AProc: TConvertFunc<T, TState>);
|
||||
end;
|
||||
|
||||
// Abstract base class for components that now act as a producer and contain a consumer.
|
||||
TMycConverter<S, T> = class abstract(TMycProducer<T>, IConverter<S, T>)
|
||||
strict private
|
||||
type
|
||||
TNull = class(TInterfacedObject, IConverter<S, T>)
|
||||
private
|
||||
function GetConsumer: IConsumer<S>;
|
||||
function Link(const Consumer: IConsumer<T>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
end;
|
||||
class var
|
||||
FNull: IConverter<S, T>;
|
||||
class constructor CreateClass;
|
||||
private
|
||||
FConsumer: TMycGenericConsumer<S>;
|
||||
function GetConsumer: IConsumer<S>;
|
||||
protected
|
||||
// To be implemented by descendants to perform the actual conversion.
|
||||
function Consume(const Value: S): TState; virtual; abstract;
|
||||
public
|
||||
constructor Create;
|
||||
destructor Destroy; override;
|
||||
|
||||
class property Null: IConverter<S, T> read FNull;
|
||||
end;
|
||||
|
||||
// Concrete producer that manages a list of consumers (listeners).
|
||||
TMycContainedProducer<T> = class(TContainedObject, IProducer<T>)
|
||||
private
|
||||
FListeners: TMycNotifyList<IConsumer<T>>;
|
||||
public
|
||||
constructor Create(const Controller: IInterface);
|
||||
destructor Destroy; override;
|
||||
function Broadcast(const Value: T): TState;
|
||||
function Link(const Consumer: IConsumer<T>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
end;
|
||||
|
||||
// A generic converter that uses a function reference for the conversion logic.
|
||||
TMycGenericConverter<S, T> = class(TMycConverter<S, T>)
|
||||
private
|
||||
FFunc: TConvertFunc<S, T>;
|
||||
protected
|
||||
function Consume(const Value: S): TState; override;
|
||||
public
|
||||
constructor Create(const AFunc: TConvertFunc<S, T>);
|
||||
end;
|
||||
|
||||
TMycGenericAggregator<S, T> = class(TMycConverter<S, T>)
|
||||
private
|
||||
FFunc: TAggregateFunc<S, T>;
|
||||
protected
|
||||
function Consume(const Value: S): TState; override;
|
||||
public
|
||||
constructor Create(const AFunc: TAggregateFunc<S, T>);
|
||||
end;
|
||||
|
||||
TMycGenericParallelConverter<S, T> = class(TMycConverter<S, T>)
|
||||
private
|
||||
FFunc: TConvertFunc<S, T>;
|
||||
FQueue: TState;
|
||||
protected
|
||||
function Consume(const Value: S): TState; override;
|
||||
public
|
||||
constructor Create(const AFunc: TConvertFunc<S, T>);
|
||||
end;
|
||||
|
||||
TMycIdentityConverter<T> = class(TMycConverter<T, T>)
|
||||
protected
|
||||
function Consume(const Value: T): TState; override; final;
|
||||
end;
|
||||
|
||||
// A converter that counts incoming data points and outputs the current count.
|
||||
TMycDataCounter<T> = class(TMycConverter<T, Int64>)
|
||||
private
|
||||
FCount: Int64;
|
||||
protected
|
||||
function Consume(const Value: T): TState; override;
|
||||
public
|
||||
constructor Create;
|
||||
end;
|
||||
|
||||
// A converter that takes an array and broadcasts each element individually.
|
||||
TMycTicker<T> = class(TMycConverter<TArray<T>, T>)
|
||||
protected
|
||||
function Consume(const Values: TArray<T>): TState; override;
|
||||
end;
|
||||
|
||||
// A converter that reads a specific field from a record using RTTI.
|
||||
TMycRecordFieldReader<S, T> = class(TMycConverter<S, T>)
|
||||
private
|
||||
FOffset: Integer;
|
||||
protected
|
||||
function Consume(const Values: S): TState; override;
|
||||
public
|
||||
constructor Create(const AFieldName: String);
|
||||
end;
|
||||
|
||||
// Endpoint that collects data into a series.
|
||||
TMycDataEndpoint<T> = class(TInterfacedObject, TLazy<TSeries<T>>.ILazy)
|
||||
private
|
||||
type
|
||||
PItem = ^TItem;
|
||||
TItem = record
|
||||
Next: PItem;
|
||||
Value: T;
|
||||
end;
|
||||
private
|
||||
FConsumer: TMycGenericConsumer<T>;
|
||||
FLookback: Integer;
|
||||
FChanged: TFlag;
|
||||
FLock: TLightweightMREW;
|
||||
FFirst: PItem;
|
||||
FCount: Integer;
|
||||
function GetChanged: TState;
|
||||
function Consume(const Value: T): TState;
|
||||
public
|
||||
constructor Create(ALookback: Integer);
|
||||
destructor Destroy; override;
|
||||
function Update(var Value: TSeries<T>): Boolean;
|
||||
|
||||
class function CreateDataEndpoint(Lookback: Integer; out Series: TLazy<TSeries<T>>): IConsumer<T>;
|
||||
end;
|
||||
|
||||
TMycParallelConverter<T> = class(TMycConverter<T, T>)
|
||||
private
|
||||
FQueue: TState;
|
||||
protected
|
||||
function Consume(const Value: T): TState; override; final;
|
||||
end;
|
||||
|
||||
IDataJoin<T> = interface(IProducer<TArray<T>>)
|
||||
function GetConsumers(Idx: Integer): IConsumer<T>;
|
||||
property Consumers[Idx: Integer]: IConsumer<T> read GetConsumers;
|
||||
end;
|
||||
|
||||
TMycDataJoinAll<T> = class(TInterfacedObject, IProducer<TArray<T>>, IDataJoin<T>)
|
||||
private
|
||||
FConsumers: array of record
|
||||
Consumer: TMycGenericConsumer<T>;
|
||||
Queue: TQueue<T>;
|
||||
end;
|
||||
|
||||
FContainedProvider: TMycContainedProducer<TArray<T>>;
|
||||
FLock: TSpinLock;
|
||||
FCount: Integer;
|
||||
function Consume(Idx: Integer; const Value: T): TState;
|
||||
function GetConsumers(Idx: Integer): IConsumer<T>;
|
||||
function Link(const Consumer: IConsumer<TArray<T>>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
public
|
||||
constructor Create(ACount: Integer);
|
||||
destructor Destroy; override;
|
||||
property Consumers[Idx: Integer]: IConsumer<T> read GetConsumers;
|
||||
end;
|
||||
|
||||
TMycDataJoinAny<T> = class(TInterfacedObject, IProducer<TArray<T>>, IDataJoin<T>)
|
||||
private
|
||||
FConsumers: array of TMycGenericConsumer<T>;
|
||||
FValues: TArray<T>;
|
||||
FInit: TArray<Boolean>;
|
||||
FInitCount: Integer;
|
||||
|
||||
FContainedProvider: TMycContainedProducer<TArray<T>>;
|
||||
FLock: TSpinLock;
|
||||
function Consume(Idx: Integer; const Value: T): TState;
|
||||
function GetConsumers(Idx: Integer): IConsumer<T>;
|
||||
function Link(const Consumer: IConsumer<TArray<T>>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
public
|
||||
constructor Create(ACount: Integer);
|
||||
destructor Destroy; override;
|
||||
property Consumers[Idx: Integer]: IConsumer<T> read GetConsumers;
|
||||
end;
|
||||
|
||||
TMycComposedConverter<S, T> = class(TInterfacedObject, IConverter<S, T>)
|
||||
private
|
||||
FConsumer: IConsumer<S>;
|
||||
FProducer: TProducer<T>;
|
||||
protected
|
||||
function GetConsumer: IConsumer<S>;
|
||||
function Link(const Consumer: IConsumer<T>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
public
|
||||
constructor Create(const AConsumer: IConsumer<S>; const AProducer: TProducer<T>);
|
||||
end;
|
||||
|
||||
TMycSequence<T> = class(TInterfacedObject, IConsumer<T>)
|
||||
private
|
||||
FProducers: TArray<TMycContainedProducer<T>>;
|
||||
protected
|
||||
function Consume(const Value: T): TState;
|
||||
function ProcessProducer(Idx: Integer; const Value: T): TState;
|
||||
public
|
||||
constructor Create(ACount: Integer);
|
||||
destructor Destroy; override;
|
||||
|
||||
class function CreateSequence(Count: Integer; out Producers: TArray<IProducer<T>>): IConsumer<T>; static;
|
||||
end;
|
||||
|
||||
implementation
|
||||
|
||||
uses
|
||||
Winapi.Windows,
|
||||
System.RTTI,
|
||||
Myc.TaskManager;
|
||||
|
||||
class constructor TMycConsumer<T>.CreateClass;
|
||||
begin
|
||||
FNull := TNull.Create;
|
||||
end;
|
||||
|
||||
{ TMycContainedProducer<T> }
|
||||
|
||||
constructor TMycContainedProducer<T>.Create(const Controller: IInterface);
|
||||
begin
|
||||
inherited Create(Controller);
|
||||
end;
|
||||
|
||||
destructor TMycContainedProducer<T>.Destroy;
|
||||
begin
|
||||
FListeners.Finalize;
|
||||
inherited Destroy;
|
||||
end;
|
||||
|
||||
function TMycContainedProducer<T>.Broadcast(const Value: T): TState;
|
||||
begin
|
||||
FListeners.Lock;
|
||||
try
|
||||
var item := FListeners.First;
|
||||
if item = nil then
|
||||
exit(TState.Null);
|
||||
|
||||
var i := 1;
|
||||
while item.Next <> nil do
|
||||
begin
|
||||
inc(i);
|
||||
item := item.Next;
|
||||
end;
|
||||
|
||||
var done := TLatch.CreateLatch(i);
|
||||
while item <> nil do
|
||||
begin
|
||||
item.Receiver.Consume(Value).Signal.Subscribe(done);
|
||||
item := item.Prev;
|
||||
end;
|
||||
|
||||
Result := done.State;
|
||||
finally
|
||||
FListeners.Release;
|
||||
end;
|
||||
end;
|
||||
|
||||
function TMycContainedProducer<T>.Link(const Consumer: IConsumer<T>): TTag;
|
||||
begin
|
||||
// Add the Consumer to the notification list
|
||||
FListeners.Lock;
|
||||
try
|
||||
Result := FListeners.Advise(Consumer);
|
||||
finally
|
||||
FListeners.Release;
|
||||
end;
|
||||
end;
|
||||
|
||||
procedure TMycContainedProducer<T>.Unlink(Tag: TTag);
|
||||
begin
|
||||
FListeners.Lock;
|
||||
try
|
||||
FListeners.Unadvise(Tag);
|
||||
finally
|
||||
FListeners.Release;
|
||||
end;
|
||||
end;
|
||||
|
||||
{ TMycConverter<S, T> }
|
||||
|
||||
constructor TMycConverter<S, T>.Create;
|
||||
begin
|
||||
inherited Create;
|
||||
FConsumer := TMycGenericConsumer<S>.Create(Self, Consume);
|
||||
end;
|
||||
|
||||
{ TMycConverter<S, T> }
|
||||
|
||||
class constructor TMycConverter<S, T>.CreateClass;
|
||||
begin
|
||||
FNull := TNull.Create;
|
||||
end;
|
||||
|
||||
destructor TMycConverter<S, T>.Destroy;
|
||||
begin
|
||||
FConsumer.Free;
|
||||
inherited Destroy;
|
||||
end;
|
||||
|
||||
function TMycConverter<S, T>.GetConsumer: IConsumer<S>;
|
||||
begin
|
||||
Result := FConsumer;
|
||||
end;
|
||||
|
||||
{ TMycConverter<S, T>.TNull }
|
||||
|
||||
function TMycConverter<S, T>.TNull.GetConsumer: IConsumer<S>;
|
||||
begin
|
||||
Result := TMycConsumer<S>.Null;
|
||||
end;
|
||||
|
||||
function TMycConverter<S, T>.TNull.Link(const Consumer: IConsumer<T>): TTag;
|
||||
begin
|
||||
Result := nil;
|
||||
end;
|
||||
|
||||
procedure TMycConverter<S, T>.TNull.Unlink(Tag: TTag);
|
||||
begin
|
||||
// NOP
|
||||
end;
|
||||
|
||||
{ TMycGenericConverter<S, T> }
|
||||
|
||||
constructor TMycGenericConverter<S, T>.Create(const AFunc: TConvertFunc<S, T>);
|
||||
begin
|
||||
inherited Create;
|
||||
FFunc := AFunc;
|
||||
end;
|
||||
|
||||
function TMycGenericConverter<S, T>.Consume(const Value: S): TState;
|
||||
begin
|
||||
Result := Broadcast(FFunc(Value));
|
||||
end;
|
||||
|
||||
{ TMycDataCounter<T> }
|
||||
|
||||
constructor TMycDataCounter<T>.Create;
|
||||
begin
|
||||
inherited Create;
|
||||
FCount := 0;
|
||||
end;
|
||||
|
||||
function TMycDataCounter<T>.Consume(const Value: T): TState;
|
||||
begin
|
||||
Result := Broadcast(FCount);
|
||||
inc(FCount);
|
||||
end;
|
||||
|
||||
{ TMycTicker<T> }
|
||||
|
||||
function TMycTicker<T>.Consume(const Values: TArray<T>): TState;
|
||||
begin
|
||||
var done := TLatch.CreateLatch(Length(Values));
|
||||
|
||||
for var i := 0 to High(Values) do
|
||||
Broadcast(Values[i]).Signal.Subscribe(done);
|
||||
|
||||
Result := done.State;
|
||||
end;
|
||||
|
||||
{ TMycRecordFieldReader<S, T> }
|
||||
|
||||
constructor TMycRecordFieldReader<S, T>.Create(const AFieldName: String);
|
||||
begin
|
||||
inherited Create;
|
||||
|
||||
var Context := TRttiContext.Create;
|
||||
var Field := Context.GetType(TypeInfo(S)).GetField(AFieldName);
|
||||
var TypeT := Context.GetType(TypeInfo(T));
|
||||
|
||||
Assert(Assigned(Field), 'Field ' + AFieldName + ' not found');
|
||||
Assert(Field.FieldType.TypeKind = TypeT.TypeKind, 'Incorrect type');
|
||||
|
||||
if Assigned(Field) and (Field.FieldType.TypeKind = TypeT.TypeKind) then
|
||||
FOffset := Field.Offset
|
||||
else
|
||||
FOffset := -1;
|
||||
end;
|
||||
|
||||
function TMycRecordFieldReader<S, T>.Consume(const Values: S): TState;
|
||||
type
|
||||
PT = ^T;
|
||||
begin
|
||||
if FOffset < 0 then
|
||||
exit(TState.Null);
|
||||
|
||||
var fieldPtr := PByte(@Values);
|
||||
inc(fieldPtr, FOffset);
|
||||
Result := Broadcast(PT(fieldPtr)^);
|
||||
end;
|
||||
|
||||
{ TMycGenericConsumer<T> }
|
||||
|
||||
constructor TMycGenericConsumer<T>.Create(const Controller: IInterface; const AProc: TConvertFunc<T, TState>);
|
||||
begin
|
||||
inherited Create(Controller);
|
||||
FProc := AProc;
|
||||
end;
|
||||
|
||||
function TMycGenericConsumer<T>.Consume(const Value: T): TState;
|
||||
begin
|
||||
Result := FProc(Value);
|
||||
end;
|
||||
|
||||
{ TMycDataEndpoint<T> }
|
||||
|
||||
constructor TMycDataEndpoint<T>.Create(ALookback: Integer);
|
||||
begin
|
||||
inherited Create;
|
||||
FLookback := ALookback;
|
||||
FConsumer := TMycGenericConsumer<T>.Create(Self, Consume);
|
||||
end;
|
||||
|
||||
destructor TMycDataEndpoint<T>.Destroy;
|
||||
begin
|
||||
FConsumer.Free;
|
||||
inherited;
|
||||
end;
|
||||
|
||||
function TMycDataEndpoint<T>.GetChanged: TState;
|
||||
begin
|
||||
Result := FChanged.State
|
||||
end;
|
||||
|
||||
function TMycDataEndpoint<T>.Consume(const Value: T): TState;
|
||||
begin
|
||||
FLock.BeginWrite;
|
||||
try
|
||||
var P: PItem;
|
||||
New(P);
|
||||
P.Next := FFirst;
|
||||
FFirst := P;
|
||||
P.Value := Value;
|
||||
inc(FCount);
|
||||
FChanged.Notify;
|
||||
finally
|
||||
FLock.EndWrite;
|
||||
end;
|
||||
end;
|
||||
|
||||
class function TMycDataEndpoint<T>.CreateDataEndpoint(Lookback: Integer; out Series: TLazy<TSeries<T>>): IConsumer<T>;
|
||||
var
|
||||
endPoint: TMycDataEndpoint<T>;
|
||||
begin
|
||||
endPoint := TMycDataEndpoint<T>.Create(Lookback);
|
||||
Result := endPoint.FConsumer;
|
||||
Series := endPoint;
|
||||
end;
|
||||
|
||||
function TMycDataEndpoint<T>.Update(var Value: TSeries<T>): Boolean;
|
||||
begin
|
||||
FLock.BeginWrite;
|
||||
try
|
||||
Result := FChanged.Reset;
|
||||
if Result then
|
||||
begin
|
||||
if FCount > 0 then
|
||||
begin
|
||||
var item: PItem;
|
||||
var Arr: TArray<T>;
|
||||
|
||||
SetLength(Arr, FCount);
|
||||
while FFirst <> nil do
|
||||
begin
|
||||
item := FFirst;
|
||||
FFirst := item.Next;
|
||||
dec(FCount);
|
||||
Assert(FCount >= 0);
|
||||
Arr[FCount] := item.Value;
|
||||
Dispose(item);
|
||||
end;
|
||||
Assert(FCount = 0);
|
||||
|
||||
Value.Add(Arr, FLookback);
|
||||
end;
|
||||
end;
|
||||
finally
|
||||
FLock.EndWrite;
|
||||
end;
|
||||
end;
|
||||
|
||||
{ TMycSequence<T> }
|
||||
|
||||
constructor TMycSequence<T>.Create(ACount: Integer);
|
||||
begin
|
||||
inherited Create;
|
||||
|
||||
SetLength(FProducers, ACount);
|
||||
for var i := 0 to High(FProducers) do
|
||||
FProducers[i] := TMycContainedProducer<T>.Create(Self);
|
||||
end;
|
||||
|
||||
destructor TMycSequence<T>.Destroy;
|
||||
begin
|
||||
for var i := High(FProducers) downto 0 do
|
||||
FProducers[i].Free;
|
||||
inherited;
|
||||
end;
|
||||
|
||||
class function TMycSequence<T>.CreateSequence(Count: Integer; out Producers: TArray<IProducer<T>>): IConsumer<T>;
|
||||
var
|
||||
Seq: TMycSequence<T>;
|
||||
begin
|
||||
Seq := TMycSequence<T>.Create(Count);
|
||||
Result := Seq;
|
||||
|
||||
SetLength(Producers, Length(Seq.FProducers));
|
||||
for var i := 0 to High(Producers) do
|
||||
Producers[i] := Seq.FProducers[i];
|
||||
end;
|
||||
|
||||
function TMycSequence<T>.Consume(const Value: T): TState;
|
||||
begin
|
||||
Result := ProcessProducer(0, Value);
|
||||
end;
|
||||
|
||||
function TMycSequence<T>.ProcessProducer(Idx: Integer; const Value: T): TState;
|
||||
begin
|
||||
if Idx >= Length(FProducers) then
|
||||
exit;
|
||||
|
||||
Result := TaskManager.RunTask(FProducers[idx].Broadcast(Value), function: TState begin Result := ProcessProducer(1 + idx, Value); end);
|
||||
end;
|
||||
|
||||
function TMycIdentityConverter<T>.Consume(const Value: T): TState;
|
||||
begin
|
||||
Result := Broadcast(Value);
|
||||
end;
|
||||
|
||||
{ TMycGenericParallelConverter<S, T> }
|
||||
|
||||
constructor TMycGenericParallelConverter<S, T>.Create(const AFunc: TConvertFunc<S, T>);
|
||||
begin
|
||||
inherited Create;
|
||||
FFunc := AFunc;
|
||||
end;
|
||||
|
||||
function TMycGenericParallelConverter<S, T>.Consume(const Value: S): TState;
|
||||
begin
|
||||
var cValue := Value;
|
||||
Result := TaskManager.RunTask(FQueue, function: TState begin Result := Broadcast(FFunc(cValue)); end);
|
||||
|
||||
FQueue := Result;
|
||||
end;
|
||||
|
||||
{ TMycParallelConverter<T> }
|
||||
|
||||
function TMycParallelConverter<T>.Consume(const Value: T): TState;
|
||||
begin
|
||||
var cValue := Value;
|
||||
Result := TaskManager.RunTask(FQueue, function: TState begin Result := Broadcast(cValue); end);
|
||||
FQueue := Result;
|
||||
end;
|
||||
|
||||
{ TMycDataJoinAll<T> }
|
||||
|
||||
constructor TMycDataJoinAll<T>.Create(ACount: Integer);
|
||||
begin
|
||||
inherited Create;
|
||||
|
||||
SetLength(FConsumers, ACount);
|
||||
|
||||
FLock := TSpinLock.Create(false);
|
||||
FCount := Length(FConsumers);
|
||||
|
||||
FContainedProvider := TMycContainedProducer<TArray<T>>.Create(Self);
|
||||
|
||||
var cFunc :=
|
||||
function(Idx: Integer): TConvertFunc<T, TState>
|
||||
begin
|
||||
Result := function(const Value: T): TState begin Result := Consume(Idx, Value); end
|
||||
end;
|
||||
|
||||
for var i := 0 to High(FConsumers) do
|
||||
with FConsumers[i] do
|
||||
begin
|
||||
Queue := TQueue<T>.Create;
|
||||
Consumer := TMycGenericConsumer<T>.Create(Self, cFunc(i));
|
||||
end;
|
||||
end;
|
||||
|
||||
destructor TMycDataJoinAll<T>.Destroy;
|
||||
begin
|
||||
for var i := High(FConsumers) downto 0 do
|
||||
with FConsumers[i] do
|
||||
begin
|
||||
Consumer.Free;
|
||||
Queue.Free;
|
||||
end;
|
||||
|
||||
FContainedProvider.Free;
|
||||
inherited;
|
||||
end;
|
||||
|
||||
function TMycDataJoinAll<T>.Link(const Consumer: IConsumer<TArray<T>>): TTag;
|
||||
begin
|
||||
Result := FContainedProvider.Link(Consumer);
|
||||
end;
|
||||
|
||||
procedure TMycDataJoinAll<T>.Unlink(Tag: TTag);
|
||||
begin
|
||||
FContainedProvider.Unlink(Tag);
|
||||
end;
|
||||
|
||||
function TMycDataJoinAll<T>.Consume(Idx: Integer; const Value: T): TState;
|
||||
begin
|
||||
var Arr: TArray<T>;
|
||||
|
||||
FLock.Enter;
|
||||
try
|
||||
if FConsumers[Idx].Queue.Count = 0 then
|
||||
dec(FCount);
|
||||
|
||||
if FCount = 0 then
|
||||
begin
|
||||
SetLength(Arr, Length(FConsumers));
|
||||
Arr[Idx] := Value;
|
||||
FCount := Length(FConsumers);
|
||||
for var i := 0 to High(FConsumers) do
|
||||
if i <> Idx then
|
||||
begin
|
||||
Arr[i] := FConsumers[i].Queue.Dequeue;
|
||||
if FConsumers[i].Queue.Count > 0 then
|
||||
dec(FCount);
|
||||
end;
|
||||
end
|
||||
else
|
||||
FConsumers[Idx].Queue.Enqueue(Value);
|
||||
finally
|
||||
FLock.Exit;
|
||||
end;
|
||||
|
||||
if Arr <> nil then
|
||||
FContainedProvider.Broadcast(Arr);
|
||||
end;
|
||||
|
||||
function TMycDataJoinAll<T>.GetConsumers(Idx: Integer): IConsumer<T>;
|
||||
begin
|
||||
Result := FConsumers[Idx].Consumer;
|
||||
end;
|
||||
|
||||
{ TMycGenericAggregator<S, T> }
|
||||
|
||||
constructor TMycGenericAggregator<S, T>.Create(const AFunc: TAggregateFunc<S, T>);
|
||||
begin
|
||||
inherited Create;
|
||||
FFunc := AFunc;
|
||||
end;
|
||||
|
||||
function TMycGenericAggregator<S, T>.Consume(const Value: S): TState;
|
||||
begin
|
||||
Result := FFunc(Value, Broadcast);
|
||||
end;
|
||||
|
||||
{ TMycComposedConverter<S, T> }
|
||||
|
||||
constructor TMycComposedConverter<S, T>.Create(const AConsumer: IConsumer<S>; const AProducer: TProducer<T>);
|
||||
begin
|
||||
inherited Create;
|
||||
FConsumer := AConsumer;
|
||||
FProducer := TProducer<T>(AProducer);
|
||||
end;
|
||||
|
||||
function TMycComposedConverter<S, T>.GetConsumer: IConsumer<S>;
|
||||
begin
|
||||
Result := FConsumer;
|
||||
end;
|
||||
|
||||
function TMycComposedConverter<S, T>.Link(const Consumer: IConsumer<T>): TTag;
|
||||
begin
|
||||
// Delegate to the contained producer
|
||||
Result := IProducer<T>(FProducer).Link(Consumer);
|
||||
end;
|
||||
|
||||
procedure TMycComposedConverter<S, T>.Unlink(Tag: TTag);
|
||||
begin
|
||||
// Delegate to the contained producer
|
||||
IProducer<T>(FProducer).Unlink(Tag);
|
||||
end;
|
||||
|
||||
{ TMycProducer<T> }
|
||||
|
||||
constructor TMycProducer<T>.Create;
|
||||
begin
|
||||
inherited Create;
|
||||
end;
|
||||
|
||||
class constructor TMycProducer<T>.CreateClass;
|
||||
begin
|
||||
FNull := TNull.Create;
|
||||
end;
|
||||
|
||||
destructor TMycProducer<T>.Destroy;
|
||||
begin
|
||||
FListeners.Finalize;
|
||||
inherited Destroy;
|
||||
end;
|
||||
|
||||
function TMycProducer<T>.Broadcast(const Value: T): TState;
|
||||
begin
|
||||
FListeners.Lock;
|
||||
try
|
||||
var item := FListeners.First;
|
||||
if item = nil then
|
||||
exit(TState.Null);
|
||||
|
||||
var i := 1;
|
||||
while item.Next <> nil do
|
||||
begin
|
||||
inc(i);
|
||||
item := item.Next;
|
||||
end;
|
||||
|
||||
var done := TLatch.CreateLatch(i);
|
||||
while item <> nil do
|
||||
begin
|
||||
item.Receiver.Consume(Value).Signal.Subscribe(done);
|
||||
item := item.Prev;
|
||||
end;
|
||||
|
||||
Result := done.State;
|
||||
finally
|
||||
FListeners.Release;
|
||||
end;
|
||||
end;
|
||||
|
||||
function TMycProducer<T>.Link(const Consumer: IConsumer<T>): TTag;
|
||||
begin
|
||||
// Add the Consumer to the notification list
|
||||
FListeners.Lock;
|
||||
try
|
||||
Result := FListeners.Advise(Consumer);
|
||||
finally
|
||||
FListeners.Release;
|
||||
end;
|
||||
end;
|
||||
|
||||
procedure TMycProducer<T>.Unlink(Tag: TTag);
|
||||
begin
|
||||
FListeners.Lock;
|
||||
try
|
||||
FListeners.Unadvise(Tag);
|
||||
finally
|
||||
FListeners.Release;
|
||||
end;
|
||||
end;
|
||||
|
||||
function TMycConsumer<T>.TNull.Consume(const Value: T): TState;
|
||||
begin
|
||||
Result := TState.Null;
|
||||
end;
|
||||
|
||||
function TMycProducer<T>.TNull.Link(const Consumer: IConsumer<T>): TTag;
|
||||
begin
|
||||
Result := nil;
|
||||
end;
|
||||
|
||||
procedure TMycProducer<T>.TNull.Unlink(Tag: TTag);
|
||||
begin
|
||||
|
||||
end;
|
||||
|
||||
{ TMycDataJoinAny<T> }
|
||||
|
||||
constructor TMycDataJoinAny<T>.Create(ACount: Integer);
|
||||
begin
|
||||
inherited Create;
|
||||
|
||||
SetLength(FConsumers, ACount);
|
||||
SetLength(FValues, ACount);
|
||||
SetLength(FInit, ACount);
|
||||
FInitCount := ACount;
|
||||
|
||||
FLock := TSpinLock.Create(false);
|
||||
|
||||
FContainedProvider := TMycContainedProducer<TArray<T>>.Create(Self);
|
||||
|
||||
var cFunc :=
|
||||
function(Idx: Integer): TConvertFunc<T, TState>
|
||||
begin
|
||||
Result := function(const Value: T): TState begin Result := Consume(Idx, Value); end
|
||||
end;
|
||||
|
||||
for var i := 0 to High(FConsumers) do
|
||||
begin
|
||||
FValues[i] := Default(T);
|
||||
FInit[i] := false;
|
||||
FConsumers[i] := TMycGenericConsumer<T>.Create(Self, cFunc(i));
|
||||
end;
|
||||
end;
|
||||
|
||||
destructor TMycDataJoinAny<T>.Destroy;
|
||||
begin
|
||||
for var i := High(FConsumers) downto 0 do
|
||||
FConsumers[i].Free;
|
||||
|
||||
FContainedProvider.Free;
|
||||
inherited;
|
||||
end;
|
||||
|
||||
function TMycDataJoinAny<T>.Link(const Consumer: IConsumer<TArray<T>>): TTag;
|
||||
begin
|
||||
Result := FContainedProvider.Link(Consumer);
|
||||
end;
|
||||
|
||||
procedure TMycDataJoinAny<T>.Unlink(Tag: TTag);
|
||||
begin
|
||||
FContainedProvider.Unlink(Tag);
|
||||
end;
|
||||
|
||||
function TMycDataJoinAny<T>.Consume(Idx: Integer; const Value: T): TState;
|
||||
begin
|
||||
FLock.Enter;
|
||||
try
|
||||
FValues[Idx] := Value;
|
||||
if not FInit[Idx] then
|
||||
begin
|
||||
FInit[Idx] := true;
|
||||
dec(FInitCount);
|
||||
end;
|
||||
|
||||
if FInitCount = 0 then
|
||||
FContainedProvider.Broadcast(FValues);
|
||||
finally
|
||||
FLock.Exit;
|
||||
end;
|
||||
end;
|
||||
|
||||
function TMycDataJoinAny<T>.GetConsumers(Idx: Integer): IConsumer<T>;
|
||||
begin
|
||||
Result := FConsumers[Idx];
|
||||
end;
|
||||
|
||||
end.
|
||||
@@ -2,8 +2,6 @@ unit Myc.Data.Pipeline;
|
||||
|
||||
interface
|
||||
|
||||
{$M+}
|
||||
|
||||
uses
|
||||
Myc.Signals,
|
||||
Myc.Mutable,
|
||||
@@ -13,12 +11,12 @@ type
|
||||
TTag = Pointer;
|
||||
|
||||
// A generic interface for components that consume data of type T.
|
||||
IConsumer<T> = interface
|
||||
IConsumer<T> = interface(IInvokable)
|
||||
function Consume(const Value: T): TState;
|
||||
end;
|
||||
|
||||
// A producer generates data and distributes it to linked consumers.
|
||||
IProducer<T> = interface
|
||||
IProducer<T> = interface(IInvokable)
|
||||
function Link(const Consumer: IConsumer<T>): TTag;
|
||||
procedure Unlink(Tag: TTag);
|
||||
end;
|
||||
|
||||
@@ -179,23 +179,27 @@ type
|
||||
procedure Add(const Item: TScalar.TValue; Lookback: Int64 = -1);
|
||||
end;
|
||||
|
||||
IRecordSeries = interface
|
||||
IScalarRecordSeries = interface
|
||||
{$region 'private'}
|
||||
function GetCount: Int64;
|
||||
function GetDef: IScalarRecordDefinition;
|
||||
function GetTotalCount: Int64;
|
||||
function GetFields(const Key: IKeyword): ISeries;
|
||||
function GetTotalCount: Int64;
|
||||
{$endregion}
|
||||
procedure Add(const Item: IScalarRecord; Lookback: Int64 = -1);
|
||||
procedure AddRaw(const Data; Lookback: Int64 = -1);
|
||||
property Count: Int64 read GetCount;
|
||||
property Def: IScalarRecordDefinition read GetDef;
|
||||
property Fields[const Key: IKeyword]: ISeries read GetFields; default;
|
||||
property TotalCount: Int64 read GetTotalCount;
|
||||
end;
|
||||
|
||||
IWriteableScalarRecordSeries = interface(IScalarRecordSeries)
|
||||
{$region 'private'}
|
||||
{$endregion}
|
||||
procedure Add(const Item: IScalarRecord; Lookback: Int64 = -1);
|
||||
end;
|
||||
|
||||
// A series of scalar records, optimized for memory and access speed.
|
||||
TScalarRecordSeries = class(TInterfacedObject, IRecordSeries)
|
||||
TScalarRecordSeries = class(TInterfacedObject, IWriteableScalarRecordSeries)
|
||||
type
|
||||
TMemberSeries = class(TGenericContainedObject<TScalarRecordSeries>, ISeries)
|
||||
private
|
||||
@@ -222,10 +226,8 @@ type
|
||||
public
|
||||
constructor Create(const ADef: IScalarRecordDefinition);
|
||||
destructor Destroy; override;
|
||||
procedure Add(const Item: IScalarRecord; Lookback: Int64 = -1);
|
||||
|
||||
// Copies data directly from the Data reference (zero-copy wrapper)
|
||||
procedure AddRaw(const Data; Lookback: Int64 = -1);
|
||||
procedure Add(const Item: IScalarRecord; Lookback: Int64 = -1);
|
||||
|
||||
property Count: Int64 read GetCount;
|
||||
property Fields[const Key: IKeyword]: ISeries read GetFields; default;
|
||||
@@ -809,18 +811,6 @@ begin
|
||||
inc(FTotalCount);
|
||||
end;
|
||||
|
||||
procedure TScalarRecordSeries.AddRaw(const Data; Lookback: Int64 = -1);
|
||||
begin
|
||||
var lb := FDef.Count * Integer(Lookback);
|
||||
var P: TScalar.PValue := @Data;
|
||||
for var i := 0 to FDef.Count - 1 do
|
||||
begin
|
||||
FArray.Add(P^, lb);
|
||||
inc(P);
|
||||
end;
|
||||
inc(FTotalCount);
|
||||
end;
|
||||
|
||||
function TScalarRecordSeries.GetFields(const Key: IKeyword): ISeries;
|
||||
begin
|
||||
var elem := FDef.IndexOf(Key);
|
||||
|
||||
@@ -78,7 +78,7 @@ type
|
||||
|
||||
class function Map(const SourceSeries: ISeries; const MapperFunc: TFunc): ISeries; static;
|
||||
|
||||
class function FromRecordSeries(const AValue: IRecordSeries): TDataValue; static; inline;
|
||||
class function FromRecordSeries(const AValue: IWriteableScalarRecordSeries): TDataValue; static; inline;
|
||||
class function FromScalarRecord(const AValue: IScalarRecord): TDataValue; static; inline;
|
||||
class function FromGenericRecord(const AValue: IKeywordMapping<TDataValue>): TDataValue; static; inline;
|
||||
class function FromSeries(const AValue: ISeries): TDataValue; static; inline;
|
||||
@@ -86,7 +86,7 @@ type
|
||||
function AsScalar: TScalar; inline;
|
||||
function AsMethod: TFunc; inline;
|
||||
function AsText: String; inline;
|
||||
function AsRecordSeries: IRecordSeries; inline;
|
||||
function AsRecordSeries: IWriteableScalarRecordSeries; inline;
|
||||
function AsScalarRecord: IScalarRecord; inline;
|
||||
function AsGenericRecord: IKeywordMapping<TDataValue>; inline;
|
||||
function AsSeries: ISeries; inline;
|
||||
@@ -205,11 +205,11 @@ begin
|
||||
Result := IKeywordMapping<TDataValue>(FInterface);
|
||||
end;
|
||||
|
||||
function TDataValue.AsRecordSeries: IRecordSeries;
|
||||
function TDataValue.AsRecordSeries: IWriteableScalarRecordSeries;
|
||||
begin
|
||||
if (FKind <> vkRecordSeries) then
|
||||
raise EInvalidCast.Create('Cannot read value as RecordSeries.');
|
||||
Result := IRecordSeries(FInterface);
|
||||
Result := IWriteableScalarRecordSeries(FInterface);
|
||||
end;
|
||||
|
||||
function TDataValue.AsScalar: TScalar;
|
||||
@@ -266,7 +266,7 @@ begin
|
||||
Result.FInterface := AValue;
|
||||
end;
|
||||
|
||||
class function TDataValue.FromRecordSeries(const AValue: IRecordSeries): TDataValue;
|
||||
class function TDataValue.FromRecordSeries(const AValue: IWriteableScalarRecordSeries): TDataValue;
|
||||
begin
|
||||
Result.FKind := vkRecordSeries;
|
||||
Result.FInterface := AValue;
|
||||
|
||||
Reference in New Issue
Block a user