DataFlow refactoring

This commit is contained in:
Michael Schimmel
2025-07-23 20:14:24 +02:00
parent b623be13fa
commit 7b2446b220
6 changed files with 260 additions and 486 deletions
+2 -2
View File
@@ -272,7 +272,7 @@ begin
inherited Create(AParent);
FDataProvider := ADataProvider;
FReceiver := TConverter.CreateEndpoint<T>(FDataProvider, Owner.Lookback.Value);
FReceiver := FDataProvider.CreateEndpoint(Owner.Lookback.Value);
end;
function TChartCustomLayer<T>.GetCount: Int64;
@@ -296,7 +296,7 @@ constructor TChartXAxisLayer<T>.Create(AOwner: TMycChart; const ADataProvider: T
begin
inherited Create(AOwner);
FDataProvider := ADataProvider;
FReceiver := TConverter.CreateEndpoint<T>(FDataProvider, Owner.Lookback.Value);
FReceiver := FDataProvider.CreateEndpoint(Owner.Lookback.Value);
end;
function TChartXAxisLayer<T>.GetCaption(Idx: Int64): String;
+1 -3
View File
@@ -636,9 +636,7 @@ function TMycChart.SetXAxisCounter<T>(const DataProvider: TDataProvider<T>): TMy
begin
FXAxisSeries.Free;
var counter := TConverter.CreateCounter<T>;
DataProvider.Link(counter);
FXAxisSeries := TChartXAxisLayer<Int64>.Create(Self, counter.Sender);
FXAxisSeries := TChartXAxisLayer<Int64>.Create(Self, DataProvider.Chain<Int64>(TConverter.CreateCounter<T>));
Result := FXAxisSeries;
end;
+26 -27
View File
@@ -22,7 +22,7 @@ type
end;
// Concrete data provider that manages a list of processors (listeners).
TMycContainedDataProvider<T> = class abstract(TContainedObject, TDataProvider<T>.IDataProvider)
TMycContainedDataProvider<T> = class abstract(TContainedObject, IDataProvider<T>)
private
FListeners: TMycNotifyList<IMycProcessor<T>>;
public
@@ -31,9 +31,9 @@ type
// Notifies all linked processors.
function Broadcast(const Value: T): TState;
// Link a Processor
function Link(const Processor: IMycProcessor<T>): TDataProvider<T>.TTag;
function Link(const Processor: IMycProcessor<T>): TTag;
// Unlink a linked strategy
procedure Unlink(Tag: TDataProvider<T>.TTag);
procedure Unlink(Tag: TTag);
end;
TMycSequence<T> = class(TMycProcessor<T>, IMycDataSequence<T>)
@@ -50,17 +50,17 @@ type
end;
// Null object implementation for IDataProvider.
TNullDataProvider<T> = class(TInterfacedObject, TDataProvider<T>.IDataProvider)
TNullDataProvider<T> = class(TInterfacedObject, IDataProvider<T>)
public
function Link(const Receiver: IMycProcessor<T>): TDataProvider<T>.TTag;
procedure Unlink(Tag: TDataProvider<T>.TTag);
function Link(const Receiver: IMycProcessor<T>): TTag;
procedure Unlink(Tag: TTag);
end;
// Abstract base class for components that process data of type S and provide data of type T.
TMycConverter<S, T> = class abstract(TMycProcessor<S>, TConverter<S, T>.IConverter)
TMycConverter<S, T> = class abstract(TMycProcessor<S>, IConverter<S, T>)
private
FSender: TMycContainedDataProvider<T>;
function GetSender: TDataProvider<T>.IDataProvider;
function GetSender: IDataProvider<T>;
protected
function ProcessData(const Value: S): TState; override; abstract;
// Broadcasts the given data to all linked processors.
@@ -68,13 +68,13 @@ type
public
constructor Create;
destructor Destroy; override;
property Sender: TDataProvider<T>.IDataProvider read GetSender;
property Sender: IDataProvider<T> read GetSender;
end;
// Null object implementation for IConverter.
TNullConverter<S, T> = class(TInterfacedObject, TConverter<S, T>.IConverter)
TNullConverter<S, T> = class(TInterfacedObject, IConverter<S, T>)
private
function GetSender: TDataProvider<T>.IDataProvider;
function GetSender: IDataProvider<T>;
public
function ProcessData(const Value: S): TState;
end;
@@ -174,8 +174,8 @@ type
end;
private
FProcessor: TMycContainedProcessor<T>;
FTag: TDataProvider<T>.TTag;
FDataProvider: TDataProvider<T>;
FTag: TTag;
FDataProvider: IDataProvider<T>;
FLookback: Int64;
FChanged: TFlag;
FLock: TLightweightMREW;
@@ -225,12 +225,12 @@ type
end;
// Endpoint that collects data into a series.
TMycDataJoin<T> = class(TInterfacedObject, TDataProvider<TArray<T>>.IDataProvider)
TMycDataJoin<T> = class(TInterfacedObject, IDataProvider<TArray<T>>)
private
FReceivers: array of record
DataProvider: TDataProvider<T>;
DataProvider: IDataProvider<T>;
Processor: TMycContainedProcessor<T>;
Tag: TDataProvider<T>.TTag;
Tag: TTag;
Queue: TQueue<T>;
end;
@@ -241,16 +241,16 @@ type
public
constructor Create(const ADataProviders: TArray<TDataProvider<T>>);
destructor Destroy; override;
property Sender: TMycContainedDataProvider<TArray<T>> read FSender implements TDataProvider<TArray<T>>.IDataProvider;
property Sender: TMycContainedDataProvider<TArray<T>> read FSender implements IDataProvider<TArray<T>>;
end;
TMycComposedConverter<S, T> = class(TMycProcessor<S>, TConverter<S, T>.IConverter)
TMycComposedConverter<S, T> = class(TMycProcessor<S>, IConverter<S, T>)
private
FProcessor: IMycProcessor<S>;
FDataProvider: TDataProvider<T>;
protected
function ProcessData(const Value: S): TState; override;
function GetSender: TDataProvider<T>.IDataProvider;
function GetSender: IDataProvider<T>;
public
constructor Create(const AProcessor: IMycProcessor<S>; const ADataProvider: TDataProvider<T>);
end;
@@ -314,7 +314,7 @@ begin
end;
end;
function TMycContainedDataProvider<T>.Link(const Processor: IMycProcessor<T>): TDataProvider<T>.TTag;
function TMycContainedDataProvider<T>.Link(const Processor: IMycProcessor<T>): TTag;
begin
// Add the Processor to the notification list
FListeners.Lock;
@@ -325,7 +325,7 @@ begin
end;
end;
procedure TMycContainedDataProvider<T>.Unlink(Tag: TDataProvider<T>.TTag);
procedure TMycContainedDataProvider<T>.Unlink(Tag: TTag);
begin
FListeners.Lock;
try
@@ -337,12 +337,12 @@ end;
{ TNullDataProvider<T> }
function TNullDataProvider<T>.Link(const Receiver: IMycProcessor<T>): TDataProvider<T>.TTag;
function TNullDataProvider<T>.Link(const Receiver: IMycProcessor<T>): TTag;
begin
Result := nil;
end;
procedure TNullDataProvider<T>.Unlink(Tag: TDataProvider<T>.TTag);
procedure TNullDataProvider<T>.Unlink(Tag: TTag);
begin
// Do nothing in the null implementation.
end;
@@ -366,14 +366,14 @@ begin
Result := FSender.Broadcast(Value);
end;
function TMycConverter<S, T>.GetSender: TDataProvider<T>.IDataProvider;
function TMycConverter<S, T>.GetSender: IDataProvider<T>;
begin
Result := FSender;
end;
{ TNullConverter<S, T> }
function TNullConverter<S, T>.GetSender: TDataProvider<T>.IDataProvider;
function TNullConverter<S, T>.GetSender: IDataProvider<T>;
begin
Result := TDataProvider<T>.Null;
end;
@@ -499,7 +499,6 @@ begin
FProcessor := TMycContainedProcessor<T>.Create(Self, ProcessData);
FTag := FDataProvider.Link(FProcessor);
end;
destructor TMycDataEndpoint<T>.Destroy;
@@ -930,7 +929,7 @@ begin
FDataProvider := ADataProvider;
end;
function TMycComposedConverter<S, T>.GetSender: TDataProvider<T>.IDataProvider;
function TMycComposedConverter<S, T>.GetSender: IDataProvider<T>;
begin
Result := FDataProvider;
end;
+88 -92
View File
@@ -15,38 +15,60 @@ type
function ProcessData(const Value: T): TState;
end;
TTag = Pointer;
IDataProvider<T> = interface
function Link(const Receiver: IMycProcessor<T>): TTag;
procedure Unlink(Tag: TTag);
end;
IConverter<S, T> = interface(IMycProcessor<S>)
function GetSender: IDataProvider<T>;
property Sender: IDataProvider<T> read GetSender;
end;
// Interface helper for IDataProvider providing the null object pattern.
TDataProvider<T> = record
public
type
TTag = Pointer;
IDataProvider = interface
function Link(const Receiver: IMycProcessor<T>): TTag;
procedure Unlink(Tag: TTag);
TLink = record
private
FDataProvider: IDataProvider<T>;
FTag: TTag;
public
procedure Unlink;
property DataProvider: IDataProvider<T> read FDataProvider;
end;
strict private
class var
FNull: IDataProvider;
FNull: IDataProvider<T>;
class constructor CreateClass;
private
FDataProvider: IDataProvider;
FDataProvider: IDataProvider<T>;
public
constructor Create(const ADataProvider: IDataProvider);
constructor Create(const ADataProvider: IDataProvider<T>);
// Managed record operators
class operator Initialize(out Dest: TDataProvider<T>);
class operator Implicit(const A: IDataProvider): TDataProvider<T>; overload;
class operator Implicit(const A: TDataProvider<T>): IDataProvider; overload;
class operator Implicit(const A: IDataProvider<T>): TDataProvider<T>; overload;
class operator Implicit(const A: TDataProvider<T>): IDataProvider<T>; overload;
// Wrapper for IMycDataProvider methods
function Link(const Receiver: IMycProcessor<T>): TTag; inline;
procedure Unlink(Tag: TTag); inline;
function Link(const Receiver: IMycProcessor<T>): TLink; inline;
function Chain(const Next: IMycProcessor<T>): IMycProcessor<T>; overload; inline;
function Chain<R>(const Next: IConverter<T, R>): TDataProvider<R>; overload; inline;
function Chain<R>(const Func: TConstFunc<T, R>): TDataProvider<R>; overload; inline;
function CreateEndpoint(Lookback: Int64): TLazy<TSeries<T>>;
// Extracts the field of a record by it's name (using RTTI).
function Field<R>(const FieldName: String): TDataProvider<R>; inline;
function MakeParallel: TDataProvider<T>; inline;
// Provides access to the null object instance.
class property Null: IDataProvider read FNull;
class property Null: IDataProvider<T> read FNull;
end;
IMycDataSequence<T> = interface(IMycProcessor<T>)
@@ -60,57 +82,38 @@ type
TConverter<S, T> = record
public
type
IConverter = interface(IMycProcessor<S>)
{$region 'private'}
function GetSender: TDataProvider<T>.IDataProvider;
{$endregion}
property Sender: TDataProvider<T>.IDataProvider read GetSender;
end;
TBroadcastProc = reference to function(const Value: T): TState;
strict private
class var
FNull: IConverter;
FNull: IConverter<S, T>;
class constructor CreateClass;
private
FConverter: IConverter;
FConverter: IConverter<S, T>;
function GetSender: TDataProvider<T>; inline;
public
constructor Create(const AConverter: IConverter);
constructor Create(const AConverter: IConverter<S, T>);
// Managed record operators
class operator Initialize(out Dest: TConverter<S, T>);
class operator Implicit(const A: IConverter): TConverter<S, T>; overload;
class operator Implicit(const A: TConverter<S, T>): IConverter; overload;
class operator Implicit(const A: IConverter<S, T>): TConverter<S, T>; overload;
class operator Implicit(const A: TConverter<S, T>): IConverter<S, T>; overload;
class function Construct(const Processor: IMycProcessor<S>; const DataProvider: TDataProvider<T>): TConverter<S, T>; static;
class function CreateGeneric(const Func: TConstFunc<S, T>): TConverter<S, T>; static;
class function CreateAggregation(const Func: TConstFunc<S, TBroadcastProc, TState>): TConverter<S, T>; static;
class function CreateParallel(const Func: TConstFunc<S, T>): TConverter<S, T>; static;
function Chain<R>(const Next: TConverter<T, R>): TConverter<T, R>; overload; inline;
function Chain<R>(const Func: TConstFunc<T, R>): TConverter<T, R>; overload; inline;
function ChainParallel<R>(const Func: TConstFunc<T, R>): TConverter<T, R>; overload; inline;
function MakeParallel: TConverter<T, T>; overload; inline;
// Extracts the field of a record by it's name (using RTTI).
function Field<R>(const FieldName: String): TConverter<T, R>; overload; inline;
function Sequence(Count: Integer): IMycDataSequence<T>; overload;
function Sequence(const Items: TArray<IConverter>): IMycDataSequence<S>; overload;
function Sequence(const Items: TArray<IConverter<S, T>>): IMycDataSequence<S>; overload;
// Provides access to the null object instance.
class property Null: IConverter read FNull;
class property Null: IConverter<S, T> read FNull;
// Wrapper for IConverter.Sender
property Sender: TDataProvider<T> read GetSender;
end;
// Factory for creating specific converter instances.
TConverter = record
class function CreateEndpoint<T>(const DataProvider: TDataProvider<T>; Lookback: Int64): TLazy<TSeries<T>>; static;
class function CreateCounter<T>: TConverter<T, Int64>; static;
class function CreateTicker<T>: TConverter<TArray<T>, T>; static;
class function CreateRecordField<S, T>(const FieldName: String): TConverter<S, T>; static;
@@ -123,8 +126,6 @@ type
class function CreateSequence<T>(Count: Integer; const Parent: TDataProvider<T>): TArray<TConverter<T, T>>; overload; static;
class function Parallel<T>(Parent: TDataProvider<T>): TConverter<T, T>; static;
class function FieldToRecord<T>(const Layout: TDataRecord.TLayout; const Name: String): TConverter<T, TDataRecord>; static;
class function FieldOfRecord<T>(const Layout: TDataRecord.TLayout; const Name: String): TConverter<TDataRecord, T>; static;
@@ -163,36 +164,64 @@ begin
FNull := TNullDataProvider<T>.Create;
end;
constructor TDataProvider<T>.Create(const ADataProvider: IDataProvider);
constructor TDataProvider<T>.Create(const ADataProvider: IDataProvider<T>);
begin
FDataProvider := ADataProvider;
if not Assigned(FDataProvider) then
FDataProvider := FNull;
end;
function TDataProvider<T>.Chain(const Next: IMycProcessor<T>): IMycProcessor<T>;
begin
FDataProvider.Link(Next);
Result := Next;
end;
function TDataProvider<T>.Chain<R>(const Next: IConverter<T, R>): TDataProvider<R>;
begin
FDataProvider.Link(Next);
Result := Next.Sender;
end;
function TDataProvider<T>.Chain<R>(const Func: TConstFunc<T, R>): TDataProvider<R>;
begin
Result := Chain<R>(TMycGenericConverter<T, R>.Create(Func));
end;
function TDataProvider<T>.CreateEndpoint(Lookback: Int64): TLazy<TSeries<T>>;
begin
Result := TMycDataEndpoint<T>.Create(FDataProvider, Lookback);
end;
function TDataProvider<T>.Field<R>(const FieldName: String): TDataProvider<R>;
begin
Result := Chain<R>(TMycRecordFieldReader<T, R>.Create(FieldName));
end;
class operator TDataProvider<T>.Initialize(out Dest: TDataProvider<T>);
begin
Dest.FDataProvider := FNull;
end;
class operator TDataProvider<T>.Implicit(const A: IDataProvider): TDataProvider<T>;
class operator TDataProvider<T>.Implicit(const A: IDataProvider<T>): TDataProvider<T>;
begin
Result.Create(A);
end;
class operator TDataProvider<T>.Implicit(const A: TDataProvider<T>): IDataProvider;
class operator TDataProvider<T>.Implicit(const A: TDataProvider<T>): IDataProvider<T>;
begin
Result := A.FDataProvider;
end;
function TDataProvider<T>.Link(const Receiver: IMycProcessor<T>): TTag;
function TDataProvider<T>.Link(const Receiver: IMycProcessor<T>): TLink;
begin
Result := FDataProvider.Link(Receiver);
Result.FDataProvider := FDataProvider;
Result.FTag := FDataProvider.Link(Receiver);
end;
procedure TDataProvider<T>.Unlink(Tag: TTag);
function TDataProvider<T>.MakeParallel: TDataProvider<T>;
begin
FDataProvider.Unlink(Tag);
Result := Chain<T>(TMycParallelConverter<T>.Create as IConverter<T, T>);
end;
{ TConverter<S, T> }
@@ -202,29 +231,13 @@ begin
FNull := TNullConverter<S, T>.Create;
end;
constructor TConverter<S, T>.Create(const AConverter: IConverter);
constructor TConverter<S, T>.Create(const AConverter: IConverter<S, T>);
begin
FConverter := AConverter;
if not Assigned(FConverter) then
FConverter := FNull;
end;
function TConverter<S, T>.Chain<R>(const Next: TConverter<T, R>): TConverter<T, R>;
begin
FConverter.Sender.Link(Next);
Result := Next;
end;
function TConverter<S, T>.Chain<R>(const Func: TConstFunc<T, R>): TConverter<T, R>;
begin
Result := Chain<R>(TMycGenericConverter<T, R>.Create(Func));
end;
function TConverter<S, T>.ChainParallel<R>(const Func: TConstFunc<T, R>): TConverter<T, R>;
begin
Result := Chain<R>(TMycGenericParallelConverter<T, R>.Create(Func));
end;
class function TConverter<S, T>.Construct(const Processor: IMycProcessor<S>; const DataProvider: TDataProvider<T>): TConverter<S, T>;
begin
Result := TMycComposedConverter<S, T>.Create(Processor, DataProvider);
@@ -240,33 +253,18 @@ begin
Result := TMycGenericConverter<S, T>.Create(Func);
end;
class function TConverter<S, T>.CreateParallel(const Func: TConstFunc<S, T>): TConverter<S, T>;
begin
Result := TMycGenericParallelConverter<S, T>.Create(Func);
end;
function TConverter<S, T>.Field<R>(const FieldName: String): TConverter<T, R>;
begin
Result := Chain<R>(TConverter.CreateRecordField<T, R>(FieldName));
end;
function TConverter<S, T>.GetSender: TDataProvider<T>;
begin
Result := FConverter.Sender;
end;
function TConverter<S, T>.MakeParallel: TConverter<T, T>;
begin
Result := Chain<T>(TMycGenericParallelConverter<T, T>.Create(function(const Value: T): T begin Result := Value; end));
end;
function TConverter<S, T>.Sequence(Count: Integer): IMycDataSequence<T>;
begin
Result := TMycSequence<T>.Create(Count);
FConverter.Sender.Link(Result);
end;
function TConverter<S, T>.Sequence(const Items: TArray<IConverter>): IMycDataSequence<S>;
function TConverter<S, T>.Sequence(const Items: TArray<IConverter<S, T>>): IMycDataSequence<S>;
begin
var seq: IMycDataSequence<S> := TMycSequence<S>.Create(Length(Items));
@@ -281,12 +279,12 @@ begin
Dest.FConverter := FNull;
end;
class operator TConverter<S, T>.Implicit(const A: IConverter): TConverter<S, T>;
class operator TConverter<S, T>.Implicit(const A: IConverter<S, T>): TConverter<S, T>;
begin
Result.Create(A);
end;
class operator TConverter<S, T>.Implicit(const A: TConverter<S, T>): IConverter;
class operator TConverter<S, T>.Implicit(const A: TConverter<S, T>): IConverter<S, T>;
begin
Result := A.FConverter;
end;
@@ -316,11 +314,6 @@ begin
);
end;
class function TConverter.CreateEndpoint<T>(const DataProvider: TDataProvider<T>; Lookback: Int64): TLazy<TSeries<T>>;
begin
Result := TMycDataEndpoint<T>.Create(DataProvider, Lookback);
end;
class function TConverter.CreateIdentity<T>: TConverter<T, T>;
begin
Result := TMycIdentityConverter<T>.Create;
@@ -423,10 +416,13 @@ begin
RecProvider.Link(Result);
end;
class function TConverter.Parallel<T>(Parent: TDataProvider<T>): TConverter<T, T>;
procedure TDataProvider<T>.TLink.Unlink;
begin
Result := TMycParallelConverter<T>.Create;
Parent.Link(Result);
if FTag <> nil then
begin
FDataProvider.Unlink(FTag);
FTag := nil;
end;
end;
end.