unit Myc.Trade.DataPoint; interface uses Myc.Signals, Myc.Mutable, Myc.Trade.Types, Myc.Trade.DataArray, Myc.DataRecord; type // A generic interface for components that process data of type T. IMycProcessor = interface function ProcessData(const Value: T): TState; end; // Interface helper for IDataProvider providing the null object pattern. TDataProvider = record public type TTag = Pointer; IDataProvider = interface function Link(const Receiver: IMycProcessor): TTag; procedure Unlink(Tag: TTag); end; strict private class var FNull: IDataProvider; class constructor CreateClass; private FDataProvider: IDataProvider; public constructor Create(const ADataProvider: IDataProvider); // Managed record operators class operator Initialize(out Dest: TDataProvider); class operator Implicit(const A: IDataProvider): TDataProvider; overload; class operator Implicit(const A: TDataProvider): IDataProvider; overload; // Wrapper for IMycDataProvider methods function Link(const Receiver: IMycProcessor): TTag; inline; procedure Unlink(Tag: TTag); inline; // Provides access to the null object instance. class property Null: IDataProvider read FNull; end; IMycDataSequence = interface(IMycProcessor) function GetCount: Integer; function GetDataProvider(Idx: Integer): TDataProvider; property Count: Integer read GetCount; property DataProvider[Idx: Integer]: TDataProvider read GetDataProvider; default; end; // Interface helper for IConverter providing the null object pattern. TConverter = record public type IConverter = interface(IMycProcessor) {$region 'private'} function GetSender: TDataProvider.IDataProvider; {$endregion} property Sender: TDataProvider.IDataProvider read GetSender; end; TBroadcastProc = reference to function(const Value: T): TState; strict private class var FNull: IConverter; class constructor CreateClass; private FConverter: IConverter; function GetSender: TDataProvider; inline; public constructor Create(const AConverter: IConverter); // Managed record operators class operator Initialize(out Dest: TConverter); class operator Implicit(const A: IConverter): TConverter; overload; class operator Implicit(const A: TConverter): IConverter; overload; class function Construct(const Processor: IMycProcessor; const DataProvider: TDataProvider): TConverter; static; class function CreateGeneric(const Func: TConstFunc): TConverter; static; class function CreateAggregation(const Func: TConstFunc): TConverter; static; class function CreateParallel(const Func: TConstFunc): TConverter; static; function Chain(const Next: TConverter): TConverter; overload; inline; function Chain(const Func: TConstFunc): TConverter; overload; inline; function ChainParallel(const Func: TConstFunc): TConverter; overload; inline; function MakeParallel: TConverter; overload; inline; // Extracts the field of a record by it's name (using RTTI). function Field(const FieldName: String): TConverter; overload; inline; function Sequence(Count: Integer): IMycDataSequence; overload; function Sequence(const Items: TArray): IMycDataSequence; overload; // Provides access to the null object instance. class property Null: IConverter read FNull; // Wrapper for IConverter.Sender property Sender: TDataProvider read GetSender; end; // Factory for creating specific converter instances. TConverter = record class function CreateEndpoint(const DataProvider: TDataProvider; Lookback: Int64): TLazy>; static; class function CreateCounter: TConverter; static; class function CreateTicker: TConverter, T>; static; class function CreateRecordField(const FieldName: String): TConverter; static; class function CreateIdentity: TConverter; static; class function CreateTickAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; class function CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; class function CreateDataPointConverter(const Func: TConstFunc): TConverter, TDataPoint>; static; class function CreateSequence(Count: Integer; const Parent: TDataProvider): TArray>; overload; static; class function Parallel(Parent: TDataProvider): TConverter; static; class function FieldToRecord(const Layout: TDataRecord.TLayout; const Name: String): TConverter; static; class function FieldOfRecord(const Layout: TDataRecord.TLayout; const Name: String): TConverter; static; class function Join(const DataProviders: TArray>): TDataProvider>; static; type TRecordMapping = record Layout: TDataRecord.TLayout; Fields: array of record FromField, ToField: String end; end; class function DataMapping( const Inputs: TArray; const Output: TDataRecord.TLayout ): TConverter, TDataRecord>; static; class function JoinRecords( const TargetLayout: TDataRecord.TLayout; const Mapping: TArray; const DataProviders: TArray> ): TConverter, TDataRecord>; static; end; implementation uses Myc.Trade.DataPoint.Impl; { TDataProvider } class constructor TDataProvider.CreateClass; begin // Create the singleton null object instance. FNull := TNullDataProvider.Create; end; constructor TDataProvider.Create(const ADataProvider: IDataProvider); begin FDataProvider := ADataProvider; if not Assigned(FDataProvider) then FDataProvider := FNull; end; class operator TDataProvider.Initialize(out Dest: TDataProvider); begin Dest.FDataProvider := FNull; end; class operator TDataProvider.Implicit(const A: IDataProvider): TDataProvider; begin Result.Create(A); end; class operator TDataProvider.Implicit(const A: TDataProvider): IDataProvider; begin Result := A.FDataProvider; end; function TDataProvider.Link(const Receiver: IMycProcessor): TTag; begin Result := FDataProvider.Link(Receiver); end; procedure TDataProvider.Unlink(Tag: TTag); begin FDataProvider.Unlink(Tag); end; { TConverter } class constructor TConverter.CreateClass; begin FNull := TNullConverter.Create; end; constructor TConverter.Create(const AConverter: IConverter); begin FConverter := AConverter; if not Assigned(FConverter) then FConverter := FNull; end; function TConverter.Chain(const Next: TConverter): TConverter; begin FConverter.Sender.Link(Next); Result := Next; end; function TConverter.Chain(const Func: TConstFunc): TConverter; begin Result := Chain(TMycGenericConverter.Create(Func)); end; function TConverter.ChainParallel(const Func: TConstFunc): TConverter; begin Result := Chain(TMycGenericParallelConverter.Create(Func)); end; class function TConverter.Construct(const Processor: IMycProcessor; const DataProvider: TDataProvider): TConverter; begin Result := TMycComposedConverter.Create(Processor, DataProvider); end; class function TConverter.CreateAggregation(const Func: TConstFunc): TConverter; begin Result := TMycGenericAggregator.Create(Func); end; class function TConverter.CreateGeneric(const Func: TConstFunc): TConverter; begin Result := TMycGenericConverter.Create(Func); end; class function TConverter.CreateParallel(const Func: TConstFunc): TConverter; begin Result := TMycGenericParallelConverter.Create(Func); end; function TConverter.Field(const FieldName: String): TConverter; begin Result := Chain(TConverter.CreateRecordField(FieldName)); end; function TConverter.GetSender: TDataProvider; begin Result := FConverter.Sender; end; function TConverter.MakeParallel: TConverter; begin Result := Chain(TMycGenericParallelConverter.Create(function(const Value: T): T begin Result := Value; end)); end; function TConverter.Sequence(Count: Integer): IMycDataSequence; begin Result := TMycSequence.Create(Count); FConverter.Sender.Link(Result); end; function TConverter.Sequence(const Items: TArray): IMycDataSequence; begin var seq: IMycDataSequence := TMycSequence.Create(Length(Items)); for var i := 0 to High(Items) do seq[i].Link(Items[i]); Result := seq; end; class operator TConverter.Initialize(out Dest: TConverter); begin Dest.FConverter := FNull; end; class operator TConverter.Implicit(const A: IConverter): TConverter; begin Result.Create(A); end; class operator TConverter.Implicit(const A: TConverter): IConverter; begin Result := A.FConverter; end; class function TConverter.CreateTickAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; begin Result := TTickAggregation.Create(Timeframe); end; { TConverter } class function TConverter.CreateCounter: TConverter; begin Result := TMycDataCounter.Create; end; class function TConverter.CreateDataPointConverter(const Func: TConstFunc): TConverter, TDataPoint>; begin var cFunc: TConstFunc := Func; Result := TMycGenericConverter, TDataPoint>.Create( function(const Value: TDataPoint): TDataPoint begin Result.Time := Value.Time; Result.Data := cFunc(Value.Data); end ); end; class function TConverter.CreateEndpoint(const DataProvider: TDataProvider; Lookback: Int64): TLazy>; begin Result := TMycDataEndpoint.Create(DataProvider, Lookback); end; class function TConverter.CreateIdentity: TConverter; begin Result := TMycIdentityConverter.Create; end; class function TConverter.CreateRecordField(const FieldName: String): TConverter; begin Result := TMycRecordFieldReader.Create(FieldName); end; class function TConverter.CreateTicker: TConverter, T>; begin Result := TMycTicker.Create; end; class function TConverter.CreateSequence(Count: Integer; const Parent: TDataProvider): TArray>; begin var seq: IMycDataSequence := TMycSequence.Create(Count); SetLength(Result, Count); for var i := 0 to High(Result) do begin Result[i] := TMycIdentityConverter.Create; seq[i].Link(Result[i]); end; Parent.Link(seq); end; class function TConverter.CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; begin Result := TOhlcAggregation.Create(Timeframe); end; class function TConverter.DataMapping( const Inputs: TArray; const Output: TDataRecord.TLayout ): TConverter, TDataRecord>; var Idxs: array of array of record FromIdx, ToIdx: Integer end; begin SetLength(Idxs, Length(Inputs)); for var i := 0 to High(Idxs) do begin SetLength(Idxs[i], Length(Inputs[i].Fields)); for var j := 0 to High(Idxs[i]) do begin Idxs[i][j].FromIdx := Inputs[i].Layout.IndexOf(Inputs[i].Fields[j].FromField); Idxs[i][j].ToIdx := Output.IndexOf(Inputs[i].Fields[j].ToField); end; end; Result := TConverter, TDataRecord>.CreateGeneric( function(const Inputs: TArray): TDataRecord begin Result := TDataRecord.Create(Output); for var i := 0 to High(Idxs) do for var j := 0 to High(Idxs[i]) do Inputs[i].CopyField(Idxs[i][j].FromIdx, Result, Idxs[i][j].ToIdx); end ); end; class function TConverter.FieldOfRecord(const Layout: TDataRecord.TLayout; const Name: String): TConverter; begin var idx := Layout.IndexOf(Name); Result := TConverter.CreateGeneric(function(const Value: TDataRecord): T begin Value.GetValue(idx, Result); end); end; class function TConverter.FieldToRecord(const Layout: TDataRecord.TLayout; const Name: String): TConverter; begin var idx := Layout.IndexOf(Name); Result := TConverter.CreateGeneric( function(const Value: T): TDataRecord begin Result := TDataRecord.Create(Layout); Result.SetValue(idx, Value); end ); end; class function TConverter.Join(const DataProviders: TArray>): TDataProvider>; begin Result := TMycDataJoin.Create(DataProviders); end; class function TConverter.JoinRecords( const TargetLayout: TDataRecord.TLayout; const Mapping: TArray; const DataProviders: TArray> ): TConverter, TDataRecord>; begin var RecProvider := Join(DataProviders); Result := DataMapping(Mapping, TargetLayout); RecProvider.Link(Result); end; class function TConverter.Parallel(Parent: TDataProvider): TConverter; begin Result := TMycParallelConverter.Create; Parent.Link(Result); end; end.