unit Myc.Data.Pipeline; interface uses Myc.Signals, Myc.Mutable, Myc.Trade.Types, Myc.Data.Series, Myc.Data.Records; type TTag = Pointer; // A generic interface for components that consume data of type T. IConsumer = interface function Consume(const Value: T): TState; end; // A producer generates data and distributes it to linked consumers. IProducer = interface function Link(const Consumer: IConsumer): TTag; procedure Unlink(Tag: TTag); end; // A converter is a producer that internally uses a consumer to transform data. IConverter = interface(IProducer) {$region 'private'} function GetConsumer: IConsumer; {$endregion} property Consumer: IConsumer read GetConsumer; end; // Interface helper for IProducer providing the null object pattern nad subscriptions for linked consumers. TProducer = record public type TSubscription = record private FOwner: IProducer; FTag: TTag; public procedure Unlink; class operator Implicit(A: TSubscription): TProducer; overload; property Owner: IProducer read FOwner; end; private FProducer: IProducer; class function GetNull: IProducer; static; inline; public constructor Create(const AProducer: IProducer); // Managed record operators class operator Initialize(out Dest: TProducer); class operator Implicit(const A: IProducer): TProducer; overload; class operator Implicit(const A: TProducer): IProducer; overload; function CreateLink(const Consumer: IConsumer): TSubscription; inline; function CreateEndpoint(Lookback: Int64; out Series: TLazy>): TSubscription; function CreateSequence(Count: Integer; out Seq: TArray>): TSubscription; overload; experimental; // Chain consumers function Chain(const Next: IConsumer): IConsumer; overload; inline; function Chain(const Next: IConverter): TProducer; overload; inline; function Chain(const Func: TConstFunc): TProducer; overload; inline; // Extracts the field of a record by it's name (using RTTI). function Field(const FieldName: String): TProducer; inline; function MakeParallel: TProducer; inline; // Provides access to the null object instance. class property Null: IProducer read GetNull; end; // Interface helper for IConverter providing the null object pattern. TConverter = record public type TBroadcastProc = reference to function(const Value: T): TState; private FConverter: IConverter; function GetConsumer: IConsumer; inline; function GetProducer: TProducer; inline; class function GetNull: IConverter; static; 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 Consumer: IConsumer; const Producer: TProducer): TConverter; static; class function CreateConverter(const Func: TConstFunc): TConverter; static; class function CreateAggregation(const Func: TConstFunc): TConverter; static; property Consumer: IConsumer read GetConsumer; property Producer: TProducer read GetProducer; class property Null: IConverter read GetNull; end; TDataRecordLayoutHelper = record helper for TDataRecord.TLayout function FieldAsRecord(const Name: String): TConverter; function FieldOfRecord(const Name: String): TConverter; end; // Factory for creating specific converter instances. TConverter = record class function CreateCounter: TConverter; static; class function CreateTicker: TConverter, T>; static; class function CreateIdentity: TConverter; static; class function CreateTickAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; class function CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; class function CreateEndpoint(Lookback: Int64; out Series: TLazy>): IConsumer; static; class function Join(const Producers: TArray>): TProducer>; 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 Producers: TArray> ): TConverter, TDataRecord>; static; end; TDataRecordHelper = record helper for TArray> function JoinRecords( const TargetLayout: TDataRecord.TLayout; const Mapping: TArray ): TProducer; experimental; end; implementation uses Myc.Data.Pipeline.Impl; { TProducer } procedure TProducer.TSubscription.Unlink; begin if FTag <> nil then begin FOwner.Unlink(FTag); FTag := nil; end; end; class operator TProducer.TSubscription.Implicit(A: TSubscription): TProducer; begin Result := A.FOwner; end; constructor TProducer.Create(const AProducer: IProducer); begin FProducer := AProducer; if not Assigned(FProducer) then FProducer := Null; end; function TProducer.Chain(const Next: IConsumer): IConsumer; begin FProducer.Link(Next); Result := Next; end; function TProducer.Chain(const Next: IConverter): TProducer; begin FProducer.Link(Next.Consumer); Result := Next; end; function TProducer.Chain(const Func: TConstFunc): TProducer; begin Result := Chain(TMycGenericConverter.Create(Func)); end; function TProducer.CreateEndpoint(Lookback: Int64; out Series: TLazy>): TSubscription; begin Result := CreateLink(TMycDataEndpoint.CreateDataEndpoint(Lookback, Series)); end; function TProducer.CreateSequence(Count: Integer; out Seq: TArray>): TSubscription; begin Result := CreateLink(TMycSequence.CreateSequence(Count, Seq)); end; function TProducer.Field(const FieldName: String): TProducer; begin Result := Chain(TMycRecordFieldReader.Create(FieldName)); end; class function TProducer.GetNull: IProducer; begin Result := TMycProducer.Null; end; class operator TProducer.Initialize(out Dest: TProducer); begin Dest.FProducer := Null; end; class operator TProducer.Implicit(const A: IProducer): TProducer; begin Result.Create(A); end; class operator TProducer.Implicit(const A: TProducer): IProducer; begin Result := A.FProducer; end; function TProducer.CreateLink(const Consumer: IConsumer): TSubscription; begin Result.FOwner := FProducer; Result.FTag := FProducer.Link(Consumer); end; function TProducer.MakeParallel: TProducer; begin Result := Chain(TMycParallelConverter.Create as IConverter); end; constructor TConverter.Create(const AConverter: IConverter); begin FConverter := AConverter; if not Assigned(FConverter) then FConverter := Null; end; class function TConverter.Construct(const Consumer: IConsumer; const Producer: TProducer): TConverter; begin Result := TMycComposedConverter.Create(Consumer, Producer); end; class function TConverter.CreateAggregation(const Func: TConstFunc): TConverter; begin Result := TMycGenericAggregator.Create(Func); end; class function TConverter.CreateConverter(const Func: TConstFunc): TConverter; begin Result := TMycGenericConverter.Create(Func); end; function TConverter.GetConsumer: IConsumer; begin Result := FConverter.Consumer; end; function TConverter.GetProducer: TProducer; begin Result := FConverter; end; class function TConverter.GetNull: IConverter; begin Result := TMycConverter.Null; end; class operator TConverter.Initialize(out Dest: TConverter); begin Dest.FConverter := Null; 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.CreateEndpoint(Lookback: Int64; out Series: TLazy>): IConsumer; begin Result := TMycDataEndpoint.CreateDataEndpoint(Lookback, Series); end; class function TConverter.CreateIdentity: TConverter; begin Result := TMycIdentityConverter.Create; end; class function TConverter.CreateTicker: TConverter, T>; begin Result := TMycTicker.Create; 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>.CreateConverter( 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.Join(const Producers: TArray>): TProducer>; var joiner: TMycDataJoin; begin joiner := TMycDataJoin.Create(Length(Producers)); Result := joiner; for var i := 0 to High(Producers) do Producers[i].Chain(joiner.Consumers[i]); end; class function TConverter.JoinRecords( const TargetLayout: TDataRecord.TLayout; const Mapping: TArray; const Producers: TArray> ): TConverter, TDataRecord>; begin var RecProvider := Join(Producers); Result := DataMapping(Mapping, TargetLayout); RecProvider.Chain(Result.Consumer); end; function TDataRecordLayoutHelper.FieldAsRecord(const Name: String): TConverter; begin var layout := Self; var idx := IndexOf(Name); Result := TConverter.CreateConverter( function(const Value: T): TDataRecord begin Result := TDataRecord.Create(layout); Result.SetValue(idx, Value); end ); end; function TDataRecordLayoutHelper.FieldOfRecord(const Name: String): TConverter; begin var idx := IndexOf(Name); Result := TConverter.CreateConverter(function(const Value: TDataRecord): T begin Value.GetValue(idx, Result); end); end; function TDataRecordHelper.JoinRecords( const TargetLayout: TDataRecord.TLayout; const Mapping: TArray ): TProducer; begin var map := TConverter.DataMapping(Mapping, TargetLayout); TConverter.Join(Self).Chain(map.Consumer); Result := map.Producer; end; end.